diff --git a/catalog/glue/glue_test.go b/catalog/glue/glue_test.go index 3e7222dab..1501b31da 100644 --- a/catalog/glue/glue_test.go +++ b/catalog/glue/glue_test.go @@ -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()) @@ -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() } diff --git a/catalog/hadoop/hadoop.go b/catalog/hadoop/hadoop.go index 76b93a55a..3bdc73ef9 100644 --- a/catalog/hadoop/hadoop.go +++ b/catalog/hadoop/hadoop.go @@ -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) } } @@ -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 } @@ -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 } @@ -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) @@ -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) } @@ -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) } }() @@ -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 } @@ -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) { diff --git a/catalog/hadoop/hadoop_test.go b/catalog/hadoop/hadoop_test.go index 2007ab402..6fac318cf 100644 --- a/catalog/hadoop/hadoop_test.go +++ b/catalog/hadoop/hadoop_test.go @@ -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 } @@ -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) diff --git a/catalog/hive/hive.go b/catalog/hive/hive.go index 604d5b1d5..0db0a7acf 100644 --- a/catalog/hive/hive.go +++ b/catalog/hive/hive.go @@ -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) } diff --git a/catalog/hive/hive_test.go b/catalog/hive/hive_test.go index 85ab555e7..1d2a38b2c 100644 --- a/catalog/hive/hive_test.go +++ b/catalog/hive/hive_test.go @@ -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()) @@ -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() } diff --git a/catalog/rest/rest_integration_test.go b/catalog/rest/rest_integration_test.go index 35dd2bd5b..7c4dfffb5 100644 --- a/catalog/rest/rest_integration_test.go +++ b/catalog/rest/rest_integration_test.go @@ -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)) @@ -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() diff --git a/catalog/rest/scan_planning_integration_test.go b/catalog/rest/scan_planning_integration_test.go index 4d42c743a..d468798eb 100644 --- a/catalog/rest/scan_planning_integration_test.go +++ b/catalog/rest/scan_planning_integration_test.go @@ -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)) diff --git a/catalog/rest/scan_planning_test.go b/catalog/rest/scan_planning_test.go index 093eca2cd..6b1720b3c 100644 --- a/catalog/rest/scan_planning_test.go +++ b/catalog/rest/scan_planning_test.go @@ -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) } diff --git a/catalog/rest/vended_creds.go b/catalog/rest/vended_creds.go index d2d3dd43f..9e0a07e2a 100644 --- a/catalog/rest/vended_creds.go +++ b/catalog/rest/vended_creds.go @@ -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) { diff --git a/catalog/sql/sql.go b/catalog/sql/sql.go index 39c2d20c4..4b36f941d 100644 --- a/catalog/sql/sql.go +++ b/catalog/sql/sql.go @@ -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) } @@ -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 diff --git a/catalog/sql/sql_integration_test.go b/catalog/sql/sql_integration_test.go index 9f8a27932..78ece6998 100644 --- a/catalog/sql/sql_integration_test.go +++ b/catalog/sql/sql_integration_test.go @@ -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) diff --git a/cmd/iceberg/clean_orphan_files_test.go b/cmd/iceberg/clean_orphan_files_test.go index b30cf6b5f..46a0c77c1 100644 --- a/cmd/iceberg/clean_orphan_files_test.go +++ b/cmd/iceberg/clean_orphan_files_test.go @@ -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) } diff --git a/internal/mock_fs.go b/internal/mock_fs.go index 8e55fe085..dae8e5e5d 100644 --- a/internal/mock_fs.go +++ b/internal/mock_fs.go @@ -19,6 +19,7 @@ package internal import ( "bytes" + "context" "errors" sio "io" "io/fs" @@ -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) @@ -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) } diff --git a/io/gocloud/blobfs/blob.go b/io/gocloud/blobfs/blob.go index 260c02912..2225af39e 100644 --- a/io/gocloud/blobfs/blob.go +++ b/io/gocloud/blobfs/blob.go @@ -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) @@ -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} diff --git a/io/gocloud/blobfs/blob_test.go b/io/gocloud/blobfs/blob_test.go index cad387779..40a5cb3b5 100644 --- a/io/gocloud/blobfs/blob_test.go +++ b/io/gocloud/blobfs/blob_test.go @@ -176,7 +176,7 @@ func TestBlobFileIOOpenMissingReturnsPathError(t *testing.T) { fileIO := testBlobFileIO(context.Background(), "my-bucket", bucket) name := "s3://my-bucket/missing.parquet" - _, err := fileIO.Open(name) + _, err := fileIO.Open(context.Background(), name) require.ErrorIs(t, err, fs.ErrNotExist) var pathErr *fs.PathError @@ -193,7 +193,7 @@ func TestBlobFileIOOpenPreprocessErrorRetainsOriginalPath(t *testing.T) { fileIO := testBlobFileIO(context.Background(), "my-bucket", bucket) name := "s3://other-bucket/file.parquet" - _, err := fileIO.Open(name) + _, err := fileIO.Open(context.Background(), name) require.Error(t, err) var pathErr *fs.PathError @@ -369,7 +369,7 @@ func TestReadAtResourceCleanup(t *testing.T) { }, } - file, err := bfs.Open("test-file") + file, err := bfs.Open(context.Background(), "test-file") require.NoError(t, err) defer file.Close() @@ -562,7 +562,7 @@ func TestBlobFileIORawQueryFragmentRoundTrip(t *testing.T) { content := []byte("data") require.NoError(t, bfs.WriteFile(location, content)) - file, err := bfs.Open(location) + file, err := bfs.Open(context.Background(), location) require.NoError(t, err) defer file.Close() @@ -631,7 +631,7 @@ func TestBlobFileIORemoveMissingFileReturnsNotExist(t *testing.T) { extractor := DefaultObjectLocationExtractor("test-bucket") bfs := New(ctx, bucket, extractor) - err := bfs.Remove("s3://test-bucket/data/nonexistent.parquet") + err := bfs.Remove(context.Background(), "s3://test-bucket/data/nonexistent.parquet") require.ErrorIs(t, err, fs.ErrNotExist) } @@ -896,7 +896,7 @@ func TestBlobFileIOPreprocessErrorRetainsOriginalPath(t *testing.T) { op string run func() error }{ - {"remove", func() error { return fileIO.Remove(name) }}, + {"remove", func() error { return fileIO.Remove(context.Background(), name) }}, {"write file", func() error { return fileIO.WriteFile(name, nil) }}, {"new writer", func() error { _, err := fileIO.NewWriter(context.Background(), name, true, nil) @@ -922,7 +922,7 @@ func TestBlobFileIORemoveFallsBackToDirectoryMarker(t *testing.T) { require.NoError(t, bucket.WriteAll(ctx, "ns/tbl/", nil, nil)) fileIO := testBlobFileIO(ctx, "test-bucket", bucket) - require.NoError(t, fileIO.Remove("s3://test-bucket/ns/tbl")) + require.NoError(t, fileIO.Remove(context.Background(), "s3://test-bucket/ns/tbl")) exists, err := bucket.Exists(ctx, "ns/tbl/") require.NoError(t, err) diff --git a/io/io.go b/io/io.go index a788a6954..2e13e8c24 100644 --- a/io/io.go +++ b/io/io.go @@ -47,6 +47,11 @@ import ( type IO interface { // Open opens the named file. // + // The context bounds the open itself (e.g. object-store metadata + // lookups); implementations backed by a plain io/fs.FS ignore it since + // fs.FS.Open takes no context. Reads on the returned File are not bound + // by ctx. + // // When Open returns an error, it should be of type *PathError // with the Op field set to "open", the Path field set to name, // and the Err field describing the problem. @@ -54,13 +59,14 @@ type IO interface { // Open should reject attempts to open names that do not satisfy // fs.ValidPath(name), returning a *PathError with Err set to // ErrInvalid or ErrNotExist. - Open(name string) (File, error) + Open(ctx context.Context, name string) (File, error) - // Remove removes the named file or (empty) directory. + // Remove removes the named file or (empty) directory. The context bounds + // the delete; local filesystem implementations ignore it. // // If there is an error, it will be of type *PathError. // Implementations must be safe for concurrent use by multiple goroutines. - Remove(name string) error + Remove(ctx context.Context, name string) error } // ReadFileIO is the interface implemented by a file system that @@ -256,7 +262,9 @@ type ioFS struct { preProcessName func(string) string } -func (f ioFS) Open(name string) (File, error) { +func (f ioFS) Open(_ context.Context, name string) (File, error) { + // The wrapped io/fs.FS has no context-aware Open, so ctx cannot be + // forwarded; it only bounds context-aware implementations. name = f.processName(name) file, err := f.fsys.Open(name) if err != nil { @@ -280,7 +288,15 @@ func (f ioFS) processName(name string) string { return name } -func (f ioFS) Remove(name string) error { +func (f ioFS) Remove(ctx context.Context, name string) error { + // Prefer a context-aware Remove when the wrapped fsys offers one, and fall + // back to the plain io/fs-style Remove otherwise. + if r, ok := f.fsys.(interface { + Remove(context.Context, string) error + }); ok { + return r.Remove(ctx, f.processName(name)) + } + r, ok := f.fsys.(interface{ Remove(name string) error }) if !ok { return errMissingRemove diff --git a/io/io_test.go b/io/io_test.go index d473789c7..3d870dd96 100644 --- a/io/io_test.go +++ b/io/io_test.go @@ -18,6 +18,7 @@ package io_test import ( + "context" "io/fs" "testing" "testing/fstest" @@ -59,12 +60,12 @@ func TestFSPreProcNameTransformsAllOperations(t *testing.T) { return "/rewritten/file.txt" }) - file, err := fsys.Open("file://bucket/original.txt") + file, err := fsys.Open(context.Background(), "file://bucket/original.txt") require.NoError(t, err) require.NoError(t, file.Close()) _, err = fsys.(iceio.ReadFileIO).ReadFile("file://bucket/original.txt") require.NoError(t, err) - require.NoError(t, fsys.Remove("file://bucket/original.txt")) + require.NoError(t, fsys.Remove(context.Background(), "file://bucket/original.txt")) require.Equal(t, "rewritten/file.txt", underlying.opened) require.Equal(t, underlying.opened, underlying.read) diff --git a/io/local.go b/io/local.go index f6329ae61..d129487d1 100644 --- a/io/local.go +++ b/io/local.go @@ -18,6 +18,7 @@ package io import ( + "context" "fmt" "io/fs" "os" @@ -68,7 +69,7 @@ func localPath(name string) (string, error) { return filepath.FromSlash(path), nil } -func (LocalFS) Open(name string) (File, error) { +func (LocalFS) Open(_ context.Context, name string) (File, error) { path, err := localPath(name) if err != nil { return nil, err @@ -110,7 +111,7 @@ func (LocalFS) WriteFile(name string, content []byte) error { return os.WriteFile(filename, content, 0o644) } -func (LocalFS) Remove(name string) error { +func (LocalFS) Remove(_ context.Context, name string) error { path, err := localPath(name) if err != nil { return err diff --git a/io/local_test.go b/io/local_test.go index addf2225e..23b47abfb 100644 --- a/io/local_test.go +++ b/io/local_test.go @@ -18,6 +18,7 @@ package io import ( + "context" "io/fs" "net/url" "os" @@ -163,7 +164,7 @@ func TestLocalFSParsesEscapedFileURIsConsistently(t *testing.T) { func TestLocalFSRejectsUnsupportedFileURIAuthority(t *testing.T) { t.Parallel() - _, err := (LocalFS{}).Open("file://remotehost/tmp/metadata.json") + _, err := (LocalFS{}).Open(context.Background(), "file://remotehost/tmp/metadata.json") require.ErrorContains(t, err, `unsupported file URI authority "remotehost"`) } diff --git a/io/mem.go b/io/mem.go index e55b671cb..7ac217f99 100644 --- a/io/mem.go +++ b/io/mem.go @@ -66,7 +66,7 @@ func NewMemFS() *MemFS { return &MemFS{files: make(map[string][]byte)} } -func (m *MemFS) Open(name string) (File, error) { +func (m *MemFS) Open(_ context.Context, name string) (File, error) { m.mu.RLock() data, ok := m.files[name] m.mu.RUnlock() @@ -79,7 +79,7 @@ func (m *MemFS) Open(name string) (File, error) { return &memFile{data: cp, name: path.Base(name)}, nil } -func (m *MemFS) Remove(name string) error { +func (m *MemFS) Remove(_ context.Context, name string) error { m.mu.Lock() defer m.mu.Unlock() diff --git a/io/mem_test.go b/io/mem_test.go index f396acbaa..e165065d8 100644 --- a/io/mem_test.go +++ b/io/mem_test.go @@ -43,7 +43,7 @@ func TestMemIO_BasicOperations(t *testing.T) { err = writeIO.WriteFile("test-file.txt", testData) require.NoError(t, err) - file, err := memIO.Open("test-file.txt") + file, err := memIO.Open(context.Background(), "test-file.txt") require.NoError(t, err) defer file.Close() @@ -51,10 +51,10 @@ func TestMemIO_BasicOperations(t *testing.T) { require.NoError(t, err) assert.Equal(t, testData, content) - err = memIO.Remove("test-file.txt") + err = memIO.Remove(context.Background(), "test-file.txt") require.NoError(t, err) - _, err = memIO.Open("test-file.txt") + _, err = memIO.Open(context.Background(), "test-file.txt") assert.Error(t, err) } @@ -78,7 +78,7 @@ func TestMemIO_Create(t *testing.T) { err = writer.Close() require.NoError(t, err) - file, err := memIO.Open("created-file.txt") + file, err := memIO.Open(context.Background(), "created-file.txt") require.NoError(t, err) defer file.Close() @@ -107,7 +107,7 @@ func TestMemIO_MultipleFiles(t *testing.T) { } for name, expectedContent := range files { - file, err := memIO.Open(name) + file, err := memIO.Open(context.Background(), name) require.NoError(t, err) content, err := io.ReadAll(file) @@ -118,17 +118,17 @@ func TestMemIO_MultipleFiles(t *testing.T) { require.NoError(t, err) } - err = memIO.Remove("file2.txt") + err = memIO.Remove(context.Background(), "file2.txt") require.NoError(t, err) - _, err = memIO.Open("file2.txt") + _, err = memIO.Open(context.Background(), "file2.txt") assert.Error(t, err) - file1, err := memIO.Open("file1.txt") + file1, err := memIO.Open(context.Background(), "file1.txt") require.NoError(t, err) file1.Close() - file3, err := memIO.Open("file3.txt") + file3, err := memIO.Open(context.Background(), "file3.txt") require.NoError(t, err) file3.Close() } @@ -138,14 +138,14 @@ func TestMemIO_RemoveMissingReturnsNotExist(t *testing.T) { memIO, err := icebergio.LoadFS(ctx, map[string]string{}, "mem://bucket/") require.NoError(t, err) - require.ErrorIs(t, memIO.Remove("does-not-exist.txt"), fs.ErrNotExist) + require.ErrorIs(t, memIO.Remove(context.Background(), "does-not-exist.txt"), fs.ErrNotExist) } func TestMemIO_FileRejectsInvalidOffsets(t *testing.T) { memIO := icebergio.NewMemFS() require.NoError(t, memIO.WriteFile("file.txt", []byte("abc"))) - file, err := memIO.Open("file.txt") + file, err := memIO.Open(context.Background(), "file.txt") require.NoError(t, err) t.Cleanup(func() { require.NoError(t, file.Close()) }) @@ -242,7 +242,7 @@ func TestMemIO_WalkDirCallbackCanRemoveFiles(t *testing.T) { return nil } - return memIO.Remove(path) + return memIO.Remove(context.Background(), path) }) }() @@ -253,9 +253,9 @@ func TestMemIO_WalkDirCallbackCanRemoveFiles(t *testing.T) { t.Fatal("WalkDir deadlocked when its callback removed a file") } - _, err := memIO.Open("mem://bucket/root/1.txt") + _, err := memIO.Open(context.Background(), "mem://bucket/root/1.txt") require.ErrorIs(t, err, fs.ErrNotExist) - _, err = memIO.Open("mem://bucket/root/2.txt") + _, err = memIO.Open(context.Background(), "mem://bucket/root/2.txt") require.ErrorIs(t, err, fs.ErrNotExist) } diff --git a/manifest.go b/manifest.go index 259c96457..a6f077372 100644 --- a/manifest.go +++ b/manifest.go @@ -18,6 +18,7 @@ package iceberg import ( + "context" "encoding/json" "errors" "fmt" @@ -420,7 +421,9 @@ func (m *manifestFile) Entries(fs iceio.IO, discardDeleted bool) iter.Seq2[Manif } func (m *manifestFile) FetchEntries(fs iceio.IO, discardDeleted bool) (_ []ManifestEntry, err error) { - f, openErr := fs.Open(m.FilePath()) + // Entries reading does not yet carry a context; Open ignores it in every + // current IO backend, so Background is safe here. + f, openErr := fs.Open(context.Background(), m.FilePath()) if openErr != nil { return nil, openErr } diff --git a/manifest_projection.go b/manifest_projection.go index c49bdeeff..03ba9c18e 100644 --- a/manifest_projection.go +++ b/manifest_projection.go @@ -18,6 +18,7 @@ package iceberg import ( + "context" "errors" "fmt" "iter" @@ -83,7 +84,9 @@ func manifestEntries( projection *ManifestEntryProjection, ) iter.Seq2[ManifestEntry, error] { return func(yield func(ManifestEntry, error) bool) { - f, err := fs.Open(m.FilePath()) + // manifestEntries has no context to thread; Open ignores ctx in every + // current IO backend, so Background is safe here. + f, err := fs.Open(context.Background(), m.FilePath()) if err != nil { yield(nil, err) diff --git a/table/all_manifests_internal_test.go b/table/all_manifests_internal_test.go index 5bda050ca..b413cf270 100644 --- a/table/all_manifests_internal_test.go +++ b/table/all_manifests_internal_test.go @@ -134,8 +134,8 @@ type manifestTrackingIO struct { delay time.Duration } -func (fs *manifestTrackingIO) Open(name string) (iceio.File, error) { - f, err := fs.IO.Open(name) +func (fs *manifestTrackingIO) Open(_ context.Context, name string) (iceio.File, error) { + f, err := fs.IO.Open(context.Background(), name) if err != nil { return nil, err } diff --git a/table/arrow_scanner_bench_test.go b/table/arrow_scanner_bench_test.go index f3a1e5aa3..7ad985ae4 100644 --- a/table/arrow_scanner_bench_test.go +++ b/table/arrow_scanner_bench_test.go @@ -576,7 +576,7 @@ func BenchmarkArrowScanTaskResidual(b *testing.B) { if err := writer.Close(); err != nil { b.Fatal(err) } - file, err := fs.Open(dataPath) + file, err := fs.Open(context.Background(), dataPath) if err != nil { b.Fatal(err) } diff --git a/table/arrow_scanner_lazy_delete_regression_test.go b/table/arrow_scanner_lazy_delete_regression_test.go index ba0dc0bf8..9cc55c679 100644 --- a/table/arrow_scanner_lazy_delete_regression_test.go +++ b/table/arrow_scanner_lazy_delete_regression_test.go @@ -49,7 +49,7 @@ type countingOpenMemFS struct { unblock <-chan struct{} } -func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { +func (f *countingOpenMemFS) Open(_ context.Context, name string) (iceio.File, error) { f.opens.Add(1) if name == f.trackedPath { f.trackedOpens.Add(1) @@ -58,7 +58,7 @@ func (f *countingOpenMemFS) Open(name string) (iceio.File, error) { <-f.unblock } - return f.MemFS.Open(name) + return f.MemFS.Open(context.Background(), name) } func newLazyDataFile(t *testing.T, path string) iceberg.DataFile { diff --git a/table/arrow_scanner_reorder_test.go b/table/arrow_scanner_reorder_test.go index a754c759c..25d589743 100644 --- a/table/arrow_scanner_reorder_test.go +++ b/table/arrow_scanner_reorder_test.go @@ -63,7 +63,7 @@ func (g *reorderGateIO) arm(headPath string, closes chan<- struct{}) { g.closed = 0 } -func (g *reorderGateIO) Open(name string) (iceio.File, error) { +func (g *reorderGateIO) Open(_ context.Context, name string) (iceio.File, error) { g.mu.Lock() headPath, gate := g.headPath, g.gate g.mu.Unlock() @@ -72,7 +72,7 @@ func (g *reorderGateIO) Open(name string) (iceio.File, error) { <-gate } - f, err := g.LocalFS.Open(name) + f, err := g.LocalFS.Open(context.Background(), name) if err != nil { return nil, err } diff --git a/table/arrow_scanner_test.go b/table/arrow_scanner_test.go index e85de0855..973021875 100644 --- a/table/arrow_scanner_test.go +++ b/table/arrow_scanner_test.go @@ -50,9 +50,9 @@ type failAfterGoodCloseFS struct { goodClosedOnce sync.Once } -func (f *failAfterGoodCloseFS) Open(name string) (iceio.File, error) { +func (f *failAfterGoodCloseFS) Open(_ context.Context, name string) (iceio.File, error) { if name == f.goodPath { - file, err := f.MemFS.Open(name) + file, err := f.MemFS.Open(context.Background(), name) if err != nil { return nil, err } @@ -79,7 +79,7 @@ func (f *failAfterGoodCloseFS) Open(name string) (iceio.File, error) { return nil, &iofs.PathError{Op: "open", Path: name, Err: iofs.ErrNotExist} } - return f.MemFS.Open(name) + return f.MemFS.Open(context.Background(), name) } type closeSignalFile struct { diff --git a/table/commit_retry_test.go b/table/commit_retry_test.go index 415276606..fe2a62f47 100644 --- a/table/commit_retry_test.go +++ b/table/commit_retry_test.go @@ -40,8 +40,13 @@ import ( // Used to verify that doCommit fails early when the FS cannot write. type readOnlyIO struct{} -func (readOnlyIO) Open(_ string) (iceio.File, error) { return nil, errors.New("readOnlyIO: no files") } -func (readOnlyIO) Remove(_ string) error { return errors.New("readOnlyIO: read-only") } +func (readOnlyIO) Open(_ context.Context, _ string) (iceio.File, error) { + return nil, errors.New("readOnlyIO: no files") +} + +func (readOnlyIO) Remove(_ context.Context, _ string) error { + return errors.New("readOnlyIO: read-only") +} // sequentialCatalog returns a predetermined error per CommitTable attempt. // If attempts exceed the len(errs) slice it returns nil (success). diff --git a/table/conflict_validation.go b/table/conflict_validation.go index 0ade50718..4cf080c2f 100644 --- a/table/conflict_validation.go +++ b/table/conflict_validation.go @@ -48,6 +48,7 @@ package table import ( "bytes" + "context" "encoding/hex" "errors" "fmt" @@ -228,12 +229,12 @@ func newConflictManifestIO(base iceio.IO) *conflictManifestIO { } } -func (c *conflictManifestIO) Open(name string) (iceio.File, error) { +func (c *conflictManifestIO) Open(ctx context.Context, name string) (iceio.File, error) { if data, ok := c.files[name]; ok { return newConflictManifestFile(name, data), nil } - f, err := c.base.Open(name) + f, err := c.base.Open(ctx, name) if err != nil { return nil, err } @@ -262,7 +263,7 @@ func (c *conflictManifestIO) Open(name string) (iceio.File, error) { }, nil } -func (c *conflictManifestIO) Remove(string) error { +func (c *conflictManifestIO) Remove(context.Context, string) error { return errConflictManifestIOReadOnly } @@ -275,7 +276,7 @@ func (c *conflictManifestIO) Stat(name string) (fs.FileInfo, error) { return statIO.Stat(name) } - f, err := c.Open(name) + f, err := c.Open(context.Background(), name) if err != nil { return nil, err } diff --git a/table/conflict_validation_bench_test.go b/table/conflict_validation_bench_test.go index 9a994fb72..c0d66a7c3 100644 --- a/table/conflict_validation_bench_test.go +++ b/table/conflict_validation_bench_test.go @@ -19,6 +19,7 @@ package table import ( "bytes" + "context" "errors" "fmt" "testing" @@ -151,10 +152,10 @@ type conflictValidationBenchmarkIO struct { bytes int64 } -func (f *conflictValidationBenchmarkIO) Open(name string) (iceio.File, error) { +func (f *conflictValidationBenchmarkIO) Open(_ context.Context, name string) (iceio.File, error) { f.opens++ - file, err := f.IO.Open(name) + file, err := f.IO.Open(context.Background(), name) if err != nil { return nil, err } diff --git a/table/conflict_validation_test.go b/table/conflict_validation_test.go index 9a968a1f8..ee5a6a0ca 100644 --- a/table/conflict_validation_test.go +++ b/table/conflict_validation_test.go @@ -19,6 +19,7 @@ package table import ( "bytes" + "context" "errors" "fmt" "io/fs" @@ -992,7 +993,7 @@ func TestConflictManifestIODoesNotCachePartialRead(t *testing.T) { tio.files[manifestPath] = []byte("abc") cache := newConflictManifestIO(&conflictValidationStatIO{trackingCallsIO: tio}) - f, err := cache.Open(manifestPath) + f, err := cache.Open(context.Background(), manifestPath) require.NoError(t, err) buf := make([]byte, 1) _, err = f.Read(buf) @@ -1006,7 +1007,7 @@ func TestConflictManifestIOIsReadOnly(t *testing.T) { tio := newTrackingCallsIO() cache := newConflictManifestIO(tio) - err := cache.Remove("manifest.avro") + err := cache.Remove(context.Background(), "manifest.avro") require.Error(t, err) assert.Zero(t, tio.removeCount["manifest.avro"]) } @@ -1021,7 +1022,7 @@ func TestConflictManifestIOStopsCachingAtByteLimit(t *testing.T) { cachedBytes: maxConflictManifestCacheBytes - 1, } - f, err := cache.Open(manifestPath) + f, err := cache.Open(context.Background(), manifestPath) require.NoError(t, err) buf := make([]byte, 2) _, err = f.Read(buf) diff --git a/table/dv/deletion_vector.go b/table/dv/deletion_vector.go index abaa707a6..929b061b8 100644 --- a/table/dv/deletion_vector.go +++ b/table/dv/deletion_vector.go @@ -20,6 +20,7 @@ package dv import ( "bytes" "cmp" + "context" "encoding/binary" "errors" "fmt" @@ -349,7 +350,9 @@ func validateDVFile(dvFile iceberg.DataFile) error { } func openDVReader(fs iceio.IO, filePath string) (*puffin.Reader, iceio.File, error) { - f, err := fs.Open(filePath) + // Puffin DV files are opened for reading; Open ignores ctx in every current + // IO backend, so no context needs threading through the loader here. + f, err := fs.Open(context.Background(), filePath) if err != nil { return nil, nil, fmt.Errorf("open DV file %s: %w", filePath, err) } diff --git a/table/dv/deletion_vector_test.go b/table/dv/deletion_vector_test.go index bd9814e65..1d794aed4 100644 --- a/table/dv/deletion_vector_test.go +++ b/table/dv/deletion_vector_test.go @@ -19,6 +19,7 @@ package dv import ( "bytes" + "context" "encoding/binary" "encoding/json" "errors" @@ -262,8 +263,8 @@ type failingOpenFS struct { err error } -func (f failingOpenFS) Open(string) (iceio.File, error) { return nil, f.err } -func (f failingOpenFS) Remove(string) error { return nil } +func (f failingOpenFS) Open(context.Context, string) (iceio.File, error) { return nil, f.err } +func (f failingOpenFS) Remove(context.Context, string) error { return nil } type countingReadIO struct { base iceio.IO @@ -271,8 +272,8 @@ type countingReadIO struct { readCalls []readAtCall } -func (f *countingReadIO) Open(name string) (iceio.File, error) { - file, err := f.base.Open(name) +func (f *countingReadIO) Open(_ context.Context, name string) (iceio.File, error) { + file, err := f.base.Open(context.Background(), name) if err != nil { return nil, err } @@ -280,8 +281,8 @@ func (f *countingReadIO) Open(name string) (iceio.File, error) { return &countingReadFile{File: file, reads: &f.reads, readCalls: &f.readCalls}, nil } -func (f *countingReadIO) Remove(name string) error { - return f.base.Remove(name) +func (f *countingReadIO) Remove(_ context.Context, name string) error { + return f.base.Remove(context.Background(), name) } func (f *countingReadIO) reset() { diff --git a/table/dv/dv_cross_client_test.go b/table/dv/dv_cross_client_test.go index 3352dfb90..69d92ff4d 100644 --- a/table/dv/dv_cross_client_test.go +++ b/table/dv/dv_cross_client_test.go @@ -18,6 +18,7 @@ package dv import ( + "context" "path/filepath" "strconv" "testing" @@ -70,7 +71,7 @@ func TestCrossClientReadJavaMultiBlobDV(t *testing.T) { require.NoError(t, err) fs := iceio.LocalFS{} - f, err := fs.Open(abs) + f, err := fs.Open(context.Background(), abs) require.NoError(t, err) t.Cleanup(func() { f.Close() }) @@ -210,7 +211,7 @@ func loadJavaPuffinBlob(t *testing.T, fixture string, blobIdx int, referencedDat require.NoError(t, err) fs := iceio.LocalFS{} - f, err := fs.Open(abs) + f, err := fs.Open(context.Background(), abs) require.NoError(t, err) defer f.Close() diff --git a/table/dv/dv_writer_test.go b/table/dv/dv_writer_test.go index dddf09589..8d8670570 100644 --- a/table/dv/dv_writer_test.go +++ b/table/dv/dv_writer_test.go @@ -520,7 +520,7 @@ func TestDVWriterFlushUnknownSpecID(t *testing.T) { require.Error(t, err) assert.Contains(t, err.Error(), "unknown partition spec id 99") - _, openErr := fs.Open(location) + _, openErr := fs.Open(context.Background(), location) assert.ErrorIs(t, openErr, stdfs.ErrNotExist) } diff --git a/table/dv_scanner_read_test.go b/table/dv_scanner_read_test.go index c466e3d30..69d25fa77 100644 --- a/table/dv_scanner_read_test.go +++ b/table/dv_scanner_read_test.go @@ -129,8 +129,8 @@ type countingDVOpenIO struct { opens atomic.Int64 } -func (c *countingDVOpenIO) Open(name string) (iceio.File, error) { - f, err := c.LocalFS.Open(name) +func (c *countingDVOpenIO) Open(_ context.Context, name string) (iceio.File, error) { + f, err := c.LocalFS.Open(context.Background(), name) if err == nil { c.opens.Add(1) } diff --git a/table/empty_scan_task_regression_test.go b/table/empty_scan_task_regression_test.go index 87c48d90b..be3abe143 100644 --- a/table/empty_scan_task_regression_test.go +++ b/table/empty_scan_task_regression_test.go @@ -96,7 +96,7 @@ func TestScanSurvivesFullyPrunedTask(t *testing.T) { require.NoError(t, pqarrow.WriteTable(arrTbl, fo, arrTbl.NumRows(), nil, pqarrow.DefaultWriterProps())) - st, err := fs.Open(path) + st, err := fs.Open(context.Background(), path) require.NoError(t, err) defer st.Close() info, err := st.Stat() diff --git a/table/equality_delete_reader_bench_test.go b/table/equality_delete_reader_bench_test.go index 176f27c9f..20c5d2756 100644 --- a/table/equality_delete_reader_bench_test.go +++ b/table/equality_delete_reader_bench_test.go @@ -259,7 +259,7 @@ func newEqualityDeleteLoadingBenchmarkInput(b *testing.B) equalityDeleteLoadingB firstPath := "mem://benchmark-lazy-equality/delete-0000.parquet" writeEqualityDeleteParquetToMemFS(b, fs.MemFS, firstPath, `[{"id": 0}, {"id": 1}, {"id": 2}, {"id": 3}, {"id": 4}, {"id": 5}, {"id": 6}, {"id": 7}, {"id": 8}, {"id": 9}]`) - file, err := fs.MemFS.Open(firstPath) + file, err := fs.MemFS.Open(context.Background(), firstPath) if err != nil { b.Fatal(err) } diff --git a/table/equality_delete_reader_internal_test.go b/table/equality_delete_reader_internal_test.go index be0b345c0..382932cbd 100644 --- a/table/equality_delete_reader_internal_test.go +++ b/table/equality_delete_reader_internal_test.go @@ -49,9 +49,9 @@ type countingEqualityDeleteOpenFS struct { opens atomic.Int64 } -func (f *countingEqualityDeleteOpenFS) Open(name string) (iceio.File, error) { +func (f *countingEqualityDeleteOpenFS) Open(_ context.Context, name string) (iceio.File, error) { f.attempts.Add(1) - file, err := f.MemFS.Open(name) + file, err := f.MemFS.Open(context.Background(), name) if err == nil { f.opens.Add(1) } diff --git a/table/equality_delete_writer_test.go b/table/equality_delete_writer_test.go index cb0537944..c9a3df6ce 100644 --- a/table/equality_delete_writer_test.go +++ b/table/equality_delete_writer_test.go @@ -99,7 +99,7 @@ func TestWriteEqualityDeleteFiles(t *testing.T) { // Verify the file was actually written to disk fs := iceio.LocalFS{} - f, err := fs.Open(df.FilePath()) + f, err := fs.Open(context.Background(), df.FilePath()) require.NoError(t, err) require.NoError(t, f.Close()) } diff --git a/table/incremental_append_scan_test.go b/table/incremental_append_scan_test.go index e883b1484..e570b56ae 100644 --- a/table/incremental_append_scan_test.go +++ b/table/incremental_append_scan_test.go @@ -375,7 +375,7 @@ func incrementalAppendTestTable(t *testing.T) *Table { } manifestList1Path := "mem://default/table-location/metadata/snap-1.avro" writeManifestList(1, manifestList1Path, []iceberg.ManifestFile{manifest1}) - listFile, err := fs.Open(manifestList1Path) + listFile, err := fs.Open(context.Background(), manifestList1Path) require.NoError(t, err) manifestList1, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) @@ -443,7 +443,7 @@ func incrementalAppendMixedOperationTable(t *testing.T) *Table { require.NoError(t, iceberg.WriteManifestList(2, &listBuf, snapshotID, nil, &sequenceNumber, 0, manifests)) require.NoError(t, fs.WriteFile(path, listBuf.Bytes())) - listFile, err := fs.Open(path) + listFile, err := fs.Open(context.Background(), path) require.NoError(t, err) manifestList, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) @@ -581,7 +581,7 @@ func incrementalAppendExpiredExclusiveTable(t *testing.T) *Table { } manifestListBPath := "mem://default/table-location/metadata/snap-b.avro" writeManifestList(2, manifestListBPath, []iceberg.ManifestFile{manifestB}) - listFile, err := fs.Open(manifestListBPath) + listFile, err := fs.Open(context.Background(), manifestListBPath) require.NoError(t, err) manifestListB, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) diff --git a/table/incremental_changelog_scan_test.go b/table/incremental_changelog_scan_test.go index 75c960e38..90b473ef2 100644 --- a/table/incremental_changelog_scan_test.go +++ b/table/incremental_changelog_scan_test.go @@ -511,7 +511,7 @@ func incrementalChangelogTestTable(t *testing.T) *Table { &sequenceNumber, 0, manifests)) require.NoError(t, fs.WriteFile(path, buf.Bytes())) - listFile, err := fs.Open(path) + listFile, err := fs.Open(context.Background(), path) require.NoError(t, err) list, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) @@ -611,7 +611,7 @@ func incrementalChangelogManifestRewriteTable(t *testing.T) *Table { &sequenceNumber, 0, manifests)) require.NoError(t, fs.WriteFile(path, buf.Bytes())) - listFile, err := fs.Open(path) + listFile, err := fs.Open(context.Background(), path) require.NoError(t, err) list, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) @@ -711,7 +711,7 @@ func incrementalChangelogDeleteManifestTable(t *testing.T) *Table { &sequenceNumber, 0, manifests)) require.NoError(t, fs.WriteFile(path, buf.Bytes())) - listFile, err := fs.Open(path) + listFile, err := fs.Open(context.Background(), path) require.NoError(t, err) list, err := iceberg.ReadManifestList(listFile) require.NoError(t, err) @@ -774,11 +774,11 @@ type countingOpenIO struct { afterOpen func() } -func (io *countingOpenIO) Open(name string) (iceio.File, error) { +func (io *countingOpenIO) Open(_ context.Context, name string) (iceio.File, error) { io.opens.Add(1) if io.afterOpen != nil { io.afterOpen() } - return io.IO.Open(name) + return io.IO.Open(context.Background(), name) } diff --git a/table/inspect_internal_test.go b/table/inspect_internal_test.go index 24ad364e4..f066d1cf1 100644 --- a/table/inspect_internal_test.go +++ b/table/inspect_internal_test.go @@ -2201,12 +2201,12 @@ type inspectManifestListContextIO struct { rejectManifestList bool } -func (fs inspectManifestListContextIO) Open(path string) (iceio.File, error) { +func (fs inspectManifestListContextIO) Open(_ context.Context, path string) (iceio.File, error) { if fs.rejectManifestList && strings.Contains(path, "manifest-list") { return nil, errors.New("manifest list opened with uncancellable IO") } - return fs.IO.Open(path) + return fs.IO.Open(context.Background(), path) } func TestInspectFilesKeepCallerContextForManifestListRead(t *testing.T) { diff --git a/table/internal/parquet_files.go b/table/internal/parquet_files.go index c398bd9c7..3e1b47ecf 100644 --- a/table/internal/parquet_files.go +++ b/table/internal/parquet_files.go @@ -105,7 +105,7 @@ const ( type parquetFormat struct{} func (parquetFormat) Open(ctx context.Context, fs iceio.IO, path string) (_ FileReader, err error) { - inputfile, err := fs.Open(path) + inputfile, err := fs.Open(ctx, path) if err != nil { return nil, err } @@ -1260,7 +1260,8 @@ func (w *ParquetFileWriter) Close() (_ iceberg.DataFile, err error) { func (w *ParquetFileWriter) Abort() error { closeErr := w.fileCloser.Close() - removeErr := w.fs.Remove(w.info.FileName) + // Abort is a best-effort cleanup with no caller context to thread. + removeErr := w.fs.Remove(context.Background(), w.info.FileName) if errors.Is(removeErr, fs.ErrNotExist) || os.IsNotExist(removeErr) { removeErr = nil } @@ -2619,7 +2620,7 @@ func checkRowGroupBloomFilters( } func (pfs *ParquetFileSource) GetReader(ctx context.Context) (result FileReader, err error) { - pf, err := pfs.fs.Open(pfs.file.FilePath()) + pf, err := pfs.fs.Open(ctx, pfs.file.FilePath()) if err != nil { return nil, err } diff --git a/table/internal/parquet_files_test.go b/table/internal/parquet_files_test.go index 663beeeaf..376cd8572 100644 --- a/table/internal/parquet_files_test.go +++ b/table/internal/parquet_files_test.go @@ -102,11 +102,11 @@ type trackingFileSystem struct { file *trackingOpenFile } -func (f *trackingFileSystem) Open(name string) (iceio.File, error) { +func (f *trackingFileSystem) Open(_ context.Context, name string) (iceio.File, error) { return f.file, nil } -func (f *trackingFileSystem) Remove(name string) error { +func (f *trackingFileSystem) Remove(_ context.Context, name string) error { return nil } @@ -1224,7 +1224,7 @@ func TestWriteDataFileNormalizesEWKBOnDisk(t *testing.T) { }, []arrow.RecordBatch{record}) require.NoError(t, err) - output, err := fs.Open("geo.parquet") + output, err := fs.Open(context.Background(), "geo.parquet") require.NoError(t, err) data, err := io.ReadAll(output) closeErr := output.Close() @@ -1375,7 +1375,7 @@ func TestParquetFileWriterAbortRemovesFile(t *testing.T) { require.NoError(t, writer.Abort()) - _, err := memFS.Open(fileName) + _, err := memFS.Open(context.Background(), fileName) require.ErrorIs(t, err, iofs.ErrNotExist) } @@ -1979,7 +1979,7 @@ func TestParquetRowGroupTargetRotatesBeforeNextBatch(t *testing.T) { }, []arrow.RecordBatch{first, second}) require.NoError(t, err) - input, err := fsys.Open("row-groups.parquet") + input, err := fsys.Open(context.Background(), "row-groups.parquet") require.NoError(t, err) defer input.Close() diff --git a/table/lazy_deletion_vector_bench_test.go b/table/lazy_deletion_vector_bench_test.go index 43f70225d..ccdb11437 100644 --- a/table/lazy_deletion_vector_bench_test.go +++ b/table/lazy_deletion_vector_bench_test.go @@ -46,8 +46,8 @@ type countingLazyDVOpenIO struct { opens atomic.Int64 } -func (f *countingLazyDVOpenIO) Open(name string) (iceio.File, error) { - file, err := f.LocalFS.Open(name) +func (f *countingLazyDVOpenIO) Open(_ context.Context, name string) (iceio.File, error) { + file, err := f.LocalFS.Open(context.Background(), name) if err == nil { f.opens.Add(1) } diff --git a/table/lazy_deletion_vector_test.go b/table/lazy_deletion_vector_test.go index 34aaff2c8..6179822fd 100644 --- a/table/lazy_deletion_vector_test.go +++ b/table/lazy_deletion_vector_test.go @@ -130,13 +130,13 @@ type failingLazyDeletionVectorIO struct { err error } -func (f *failingLazyDeletionVectorIO) Open(string) (iceio.File, error) { +func (f *failingLazyDeletionVectorIO) Open(context.Context, string) (iceio.File, error) { f.opens.Add(1) return nil, f.err } -func (f *failingLazyDeletionVectorIO) Remove(string) error { return nil } +func (f *failingLazyDeletionVectorIO) Remove(context.Context, string) error { return nil } func TestLazyDeletionVectorLoaderCachesGroupErrors(t *testing.T) { const dataFilePath = "file:///table/data/missing.parquet" @@ -173,8 +173,8 @@ type cancelOnOpenIO struct { opens atomic.Int64 } -func (f *cancelOnOpenIO) Open(name string) (iceio.File, error) { - file, err := f.LocalFS.Open(name) +func (f *cancelOnOpenIO) Open(_ context.Context, name string) (iceio.File, error) { + file, err := f.LocalFS.Open(context.Background(), name) if err == nil { f.opens.Add(1) f.cancel() @@ -224,8 +224,8 @@ type countingPuffinOpenIO struct { opens atomic.Int64 } -func (f *countingPuffinOpenIO) Open(name string) (iceio.File, error) { - file, err := f.base.Open(name) +func (f *countingPuffinOpenIO) Open(_ context.Context, name string) (iceio.File, error) { + file, err := f.base.Open(context.Background(), name) if err == nil && strings.HasSuffix(name, ".puffin") { f.opens.Add(1) } @@ -233,7 +233,9 @@ func (f *countingPuffinOpenIO) Open(name string) (iceio.File, error) { return file, err } -func (f *countingPuffinOpenIO) Remove(name string) error { return f.base.Remove(name) } +func (f *countingPuffinOpenIO) Remove(_ context.Context, name string) error { + return f.base.Remove(context.Background(), name) +} func TestReadTasksDoesNotLoadDeletionVectorsBeforeIteration(t *testing.T) { baseFS := iceio.LocalFS{} diff --git a/table/manifest_list_partitions_regression_test.go b/table/manifest_list_partitions_regression_test.go index 998a7727c..56c3afa4a 100644 --- a/table/manifest_list_partitions_regression_test.go +++ b/table/manifest_list_partitions_regression_test.go @@ -84,7 +84,7 @@ type manifestListRecordV2 struct { func readManifestListPartitionFields(t *testing.T, fs iceio.IO, manifestListPath string) []any { t.Helper() - f, err := fs.Open(manifestListPath) + f, err := fs.Open(context.Background(), manifestListPath) require.NoError(t, err) defer func() { require.NoError(t, f.Close()) }() diff --git a/table/orphan_cleanup.go b/table/orphan_cleanup.go index dde04b57d..2719e7945 100644 --- a/table/orphan_cleanup.go +++ b/table/orphan_cleanup.go @@ -81,7 +81,7 @@ type orphanCleanupConfig struct { location string olderThan time.Duration dryRun bool - deleteFunc func(string) error + deleteFunc func(context.Context, string) error maxConcurrency int prefixMismatchMode PrefixMismatchMode equalSchemes map[string]string @@ -120,7 +120,7 @@ func WithDryRun(enabled bool) OrphanCleanupOption { // WithDeleteFunc sets a custom delete function. If not provided, the table's FileIO // delete method will be used. -func WithDeleteFunc(deleteFunc func(string) error) OrphanCleanupOption { +func WithDeleteFunc(deleteFunc func(context.Context, string) error) OrphanCleanupOption { return func(cfg *orphanCleanupConfig) { cfg.deleteFunc = deleteFunc } @@ -752,7 +752,7 @@ func deleteFilesSequential(ctx context.Context, fs iceio.IO, orphanFiles []strin break } - if err := deleteFunc(file); err != nil { + if err := deleteFunc(ctx, file); err != nil { result = errors.Join(result, fmt.Errorf("failed to delete orphan file %s: %w", file, err)) continue @@ -768,7 +768,7 @@ func deleteFilesParallel( ctx context.Context, files []string, maxConcurrency int, - deleteFunc func(string) error, + deleteFunc func(context.Context, string) error, wrapError func(string, error) error, ) ([]string, error) { workers := min(max(maxConcurrency, 1), len(files)) @@ -807,7 +807,7 @@ func deleteFilesParallel( return } - if err := deleteFunc(files[index]); err != nil { + if err := deleteFunc(ctx, files[index]); err != nil { wrappedErr := wrapError(files[index], err) if wrappedErr == nil { // Keep failed deletions observable if the wrapper suppresses the error. @@ -1401,8 +1401,8 @@ func (t Table) PurgeFiles(ctx context.Context) error { ctx, files, defaultPurgeMaxConcurrency, - func(file string) error { - if err := fs.Remove(file); err != nil && !os.IsNotExist(err) { + func(ctx context.Context, file string) error { + if err := fs.Remove(ctx, file); err != nil && !os.IsNotExist(err) { return err } diff --git a/table/orphan_cleanup_bench_test.go b/table/orphan_cleanup_bench_test.go index 8e7dd5d4c..1e190fc5d 100644 --- a/table/orphan_cleanup_bench_test.go +++ b/table/orphan_cleanup_bench_test.go @@ -115,8 +115,8 @@ type benchmarkDelayIO struct { delay time.Duration } -func (fs *benchmarkDelayIO) Open(name string) (iceio.File, error) { - f, err := fs.IO.Open(name) +func (fs *benchmarkDelayIO) Open(_ context.Context, name string) (iceio.File, error) { + f, err := fs.IO.Open(context.Background(), name) if err != nil { return nil, err } @@ -139,7 +139,7 @@ func BenchmarkPurgeFilesNonBulkDeletion(b *testing.B) { for _, delay := range []time.Duration{0, 100 * time.Microsecond, time.Millisecond} { for _, concurrency := range []int{1, 4, 16, 32, 64} { b.Run(fmt.Sprintf("files=%d/delay=%s/concurrency=%d", fileCount, delay, concurrency), func(b *testing.B) { - deleteFunc := func(string) error { + deleteFunc := func(context.Context, string) error { time.Sleep(delay) return nil diff --git a/table/orphan_cleanup_integration_test.go b/table/orphan_cleanup_integration_test.go index 0b8e0dc38..370004148 100644 --- a/table/orphan_cleanup_integration_test.go +++ b/table/orphan_cleanup_integration_test.go @@ -194,7 +194,7 @@ func (s *OrphanCleanupIntegrationSuite) mustFS(tbl *table.Table) io.IO { } func (s *OrphanCleanupIntegrationSuite) fileExists(fs io.IO, path string) bool { - f, err := fs.Open(path) + f, err := fs.Open(s.ctx, path) if err != nil { return false } @@ -248,7 +248,7 @@ func (s *OrphanCleanupIntegrationSuite) TestOrphanCleanupDryRun() { } for _, orphanFile := range orphanFiles { - s.NoError(fs.Remove(orphanFile)) + s.NoError(fs.Remove(s.ctx, orphanFile)) } } @@ -380,7 +380,7 @@ func (s *OrphanCleanupIntegrationSuite) TestOrphanCleanupWithConcurrency() { // Clean up original orphan files fs := s.mustFS(tbl) for _, orphanFile := range orphanFiles { - fs.Remove(orphanFile) // Ignore errors, might already be deleted + fs.Remove(s.ctx, orphanFile) // Ignore errors, might already be deleted } } @@ -392,10 +392,10 @@ func (s *OrphanCleanupIntegrationSuite) TestOrphanCleanupCustomDeleteFunction() time.Sleep(100 * time.Millisecond) var deletedByCustomFunc []string - customDeleteFunc := func(filePath string) error { + customDeleteFunc := func(ctx context.Context, filePath string) error { deletedByCustomFunc = append(deletedByCustomFunc, filePath) fs := s.mustFS(tbl) - return fs.Remove(filePath) + return fs.Remove(ctx, filePath) } result, err := tbl.DeleteOrphanFiles(s.ctx, diff --git a/table/orphan_cleanup_test.go b/table/orphan_cleanup_test.go index 962ff670e..bc723a637 100644 --- a/table/orphan_cleanup_test.go +++ b/table/orphan_cleanup_test.go @@ -72,7 +72,7 @@ func TestOrphanCleanupOptions(t *testing.T) { WithDryRun(true)(cfg) assert.True(t, cfg.dryRun) - deleteFunc := func(string) error { return nil } + deleteFunc := func(context.Context, string) error { return nil } WithDeleteFunc(deleteFunc)(cfg) assert.NotNil(t, cfg.deleteFunc) @@ -199,9 +199,9 @@ func TestOrphanCleanupPlanDoesNotExpandAfterPlanning(t *testing.T) { require.NoError(t, err) assert.Equal(t, []string{plannedOrphan}, result.DeletedFiles) - _, err = fs.Open(plannedOrphan) + _, err = fs.Open(context.Background(), plannedOrphan) assert.ErrorIs(t, err, stdfs.ErrNotExist) - file, err := fs.Open(newOrphan) + file, err := fs.Open(context.Background(), newOrphan) require.NoError(t, err) assert.NoError(t, file.Close()) } @@ -1515,11 +1515,11 @@ type mockBulkRemovableIO struct { bulkPaths []string } -func (m *mockBulkRemovableIO) Open(string) (io.File, error) { +func (m *mockBulkRemovableIO) Open(context.Context, string) (io.File, error) { return nil, errors.New("not implemented") } -func (m *mockBulkRemovableIO) Remove(string) error { +func (m *mockBulkRemovableIO) Remove(context.Context, string) error { return errors.New("Remove should not be called when BulkRemovableIO is available") } @@ -1547,7 +1547,7 @@ func TestDeleteFilesWithCustomDeleteFunc(t *testing.T) { var customDeleted []string cfg := &orphanCleanupConfig{ - deleteFunc: func(path string) error { + deleteFunc: func(_ context.Context, path string) error { customDeleted = append(customDeleted, path) return nil @@ -1575,11 +1575,11 @@ type mockPlainIO struct { removed []string } -func (m *mockPlainIO) Open(string) (io.File, error) { +func (m *mockPlainIO) Open(context.Context, string) (io.File, error) { return nil, errors.New("not implemented") } -func (m *mockPlainIO) Remove(name string) error { +func (m *mockPlainIO) Remove(_ context.Context, name string) error { m.mu.Lock() defer m.mu.Unlock() @@ -1635,7 +1635,7 @@ type purgeDeleteTrackingIO struct { removed []string } -func (m *purgeDeleteTrackingIO) Remove(name string) error { +func (m *purgeDeleteTrackingIO) Remove(_ context.Context, name string) error { m.mu.Lock() m.active++ m.maxActive = max(m.maxActive, m.active) @@ -1676,7 +1676,7 @@ func TestDeleteFilesParallelCollectsPurgeErrors(t *testing.T) { context.Background(), files, 4, - func(path string) error { + func(_ context.Context, path string) error { mu.Lock() calls[path]++ mu.Unlock() @@ -1730,7 +1730,7 @@ func TestDeleteFilesParallelPreservesErrorWithNilWrapper(t *testing.T) { context.Background(), files, 2, - func(path string) error { + func(_ context.Context, path string) error { if path == files[1] { return deleteErr } @@ -1753,7 +1753,7 @@ func TestDeleteFilesParallelStopsQueuedWorkOnCancellation(t *testing.T) { ctx, []string{"s3://bucket/table/file.parquet"}, 4, - func(string) error { + func(context.Context, string) error { called = true return nil @@ -1813,7 +1813,7 @@ func TestDeleteFilesParallelStopsQueuedWorkOnMidFlightCancellation(t *testing.T) gatedErrContext{Context: ctx, checks: checks, errCalls: errCalls}, files, maxConcurrency, - func(string) error { + func(context.Context, string) error { if calls.Add(1) > maxConcurrency { extraCalls <- struct{}{} } @@ -2010,7 +2010,7 @@ func TestDeleteOrphanFilesPrefixMismatchModes(t *testing.T) { WithLocation("s3://bucket/path"), WithFilesOlderThan(0), WithPrefixMismatchMode(mode), - WithDeleteFunc(func(path string) error { + WithDeleteFunc(func(_ context.Context, path string) error { deleted = append(deleted, path) return nil diff --git a/table/partition_residual_read_test.go b/table/partition_residual_read_test.go index f9bfc6021..5fb45a2ea 100644 --- a/table/partition_residual_read_test.go +++ b/table/partition_residual_read_test.go @@ -188,7 +188,7 @@ func writePartitionResidualParquet( parquet.NewWriterProperties(parquet.WithStats(true)), pqarrow.DefaultWriterProps())) require.NoError(t, writer.Close()) - file, err := fs.Open(path) + file, err := fs.Open(context.Background(), path) require.NoError(t, err) defer file.Close() info, err := file.Stat() diff --git a/table/rewrite_data_files.go b/table/rewrite_data_files.go index a66e90da2..3b58eca98 100644 --- a/table/rewrite_data_files.go +++ b/table/rewrite_data_files.go @@ -783,7 +783,7 @@ func cleanupCompactionOutputs(fs iceio.IO, batchResults []CompactionGroupResult) var cleanupErr error for path := range paths { - if err := fs.Remove(path); err != nil && !errors.Is(err, iofs.ErrNotExist) { + if err := fs.Remove(context.Background(), path); err != nil && !errors.Is(err, iofs.ErrNotExist) { cleanupErr = errors.Join(cleanupErr, fmt.Errorf("remove %s: %w", path, err)) } } diff --git a/table/rewrite_manifests.go b/table/rewrite_manifests.go index d12cbcaaf..e57ad73f8 100644 --- a/table/rewrite_manifests.go +++ b/table/rewrite_manifests.go @@ -314,7 +314,7 @@ func (r *rewriteManifests) deleteWritten(input, merged []iceberg.ManifestFile) { if _, ok := inPaths[m.FilePath()]; ok { continue } - if err := r.base.io.Remove(m.FilePath()); err != nil { + if err := r.base.io.Remove(context.Background(), m.FilePath()); err != nil { log.Printf("Warning: failed to delete orphaned merged manifest %s: %v", m.FilePath(), err) } } diff --git a/table/rewrite_manifests_bench_test.go b/table/rewrite_manifests_bench_test.go index f6d59e14c..0f8892857 100644 --- a/table/rewrite_manifests_bench_test.go +++ b/table/rewrite_manifests_bench_test.go @@ -107,7 +107,7 @@ func benchmarkManifestMergeMode(b *testing.B, cluster bool) { } b.StopTimer() for _, manifest := range output { - if err := mem.Remove(manifest.FilePath()); err != nil { + if err := mem.Remove(context.Background(), manifest.FilePath()); err != nil { b.Fatal(err) } } diff --git a/table/rewrite_manifests_cluster.go b/table/rewrite_manifests_cluster.go index 1f26f5778..fbc6b48d3 100644 --- a/table/rewrite_manifests_cluster.go +++ b/table/rewrite_manifests_cluster.go @@ -18,6 +18,7 @@ package table import ( + "context" "errors" "fmt" "io" @@ -152,7 +153,7 @@ func (m *manifestMergeManager) clusterManifests(manifests []iceberg.ManifestFile writer.abort() } for _, path := range paths { - if removeErr := m.snap.io.Remove(path); removeErr != nil { + if removeErr := m.snap.io.Remove(context.Background(), path); removeErr != nil { log.Printf("Warning: failed to delete orphaned clustered manifest %s: %v", path, removeErr) } } diff --git a/table/rewrite_manifests_test.go b/table/rewrite_manifests_test.go index b5b205c11..96f91cc25 100644 --- a/table/rewrite_manifests_test.go +++ b/table/rewrite_manifests_test.go @@ -815,12 +815,12 @@ func (f *trackingFS) Create(name string) (iceio.FileWriter, error) { return f.LocalFS.Create(name) } -func (f *trackingFS) Remove(name string) error { +func (f *trackingFS) Remove(_ context.Context, name string) error { f.mu.Lock() f.removed = append(f.removed, name) f.mu.Unlock() - return f.LocalFS.Remove(name) + return f.LocalFS.Remove(context.Background(), name) } func (f *trackingFS) snapshotCreated() map[string]struct{} { diff --git a/table/row_lineage_rewrite_bench_test.go b/table/row_lineage_rewrite_bench_test.go index b9afd326c..fb74dc542 100644 --- a/table/row_lineage_rewrite_bench_test.go +++ b/table/row_lineage_rewrite_bench_test.go @@ -158,7 +158,7 @@ func (f *rowLineageRewriteBenchmarkFixture) reset(b *testing.B) { if _, ok := f.initialFiles[path]; ok { continue } - require.NoError(b, f.fs.Remove(path)) + require.NoError(b, f.fs.Remove(context.Background(), path)) } f.catalog.metadata = f.baseMetadata } diff --git a/table/row_lineage_scan_test.go b/table/row_lineage_scan_test.go index 619c872e1..aa845df8d 100644 --- a/table/row_lineage_scan_test.go +++ b/table/row_lineage_scan_test.go @@ -171,7 +171,7 @@ func parquetRowGroupCount(t *testing.T, tbl *table.Table) int { require.NoError(t, err) require.Len(t, tasks, 1) - f, err := iceio.LocalFS{}.Open(tasks[0].File.FilePath()) + f, err := iceio.LocalFS{}.Open(context.Background(), tasks[0].File.FilePath()) require.NoError(t, err) rdr, err := file.NewParquetReader(f) require.NoError(t, err) diff --git a/table/scan_splits_test.go b/table/scan_splits_test.go index b7d9937a9..89b067027 100644 --- a/table/scan_splits_test.go +++ b/table/scan_splits_test.go @@ -472,8 +472,8 @@ type prepareReadTrackingIO struct { closed int } -func (fs *prepareReadTrackingIO) Open(path string) (iceio.File, error) { - file, err := fs.IO.Open(path) +func (fs *prepareReadTrackingIO) Open(_ context.Context, path string) (iceio.File, error) { + file, err := fs.IO.Open(context.Background(), path) if err != nil { return nil, err } diff --git a/table/scanner_internal_test.go b/table/scanner_internal_test.go index cb48bee39..e7756d844 100644 --- a/table/scanner_internal_test.go +++ b/table/scanner_internal_test.go @@ -82,12 +82,12 @@ type expiringManifestIO struct { expired atomic.Bool } -func (fs *expiringManifestIO) Open(name string) (iceio.File, error) { +func (fs *expiringManifestIO) Open(_ context.Context, name string) (iceio.File, error) { if name != fs.manifestList && fs.expired.Load() { return nil, errors.New("credentials expired") } - file, err := fs.base.Open(name) + file, err := fs.base.Open(context.Background(), name) if err == nil && name == fs.manifestList { fs.expired.Store(true) } @@ -95,8 +95,8 @@ func (fs *expiringManifestIO) Open(name string) (iceio.File, error) { return file, err } -func (fs *expiringManifestIO) Remove(name string) error { - return fs.base.Remove(name) +func (fs *expiringManifestIO) Remove(_ context.Context, name string) error { + return fs.base.Remove(context.Background(), name) } type failingManifestIO struct { @@ -105,18 +105,18 @@ type failingManifestIO struct { laterOpens atomic.Int64 } -func (fs *failingManifestIO) Open(name string) (iceio.File, error) { +func (fs *failingManifestIO) Open(_ context.Context, name string) (iceio.File, error) { if name == fs.failingPath { return nil, errors.New("manifest failed") } fs.laterOpens.Add(1) - return fs.base.Open(name) + return fs.base.Open(context.Background(), name) } -func (fs *failingManifestIO) Remove(name string) error { - return fs.base.Remove(name) +func (fs *failingManifestIO) Remove(_ context.Context, name string) error { + return fs.base.Remove(context.Background(), name) } func TestPlanFilesRefreshesFileIODuringManifestPlanning(t *testing.T) { diff --git a/table/snapshot_manifest_cache_bench_test.go b/table/snapshot_manifest_cache_bench_test.go index 1624d8c4c..c9d5507f9 100644 --- a/table/snapshot_manifest_cache_bench_test.go +++ b/table/snapshot_manifest_cache_bench_test.go @@ -89,10 +89,10 @@ type snapshotManifestCacheBenchmarkIO struct { opens map[string]int } -func (fs *snapshotManifestCacheBenchmarkIO) Open(name string) (iceio.File, error) { +func (fs *snapshotManifestCacheBenchmarkIO) Open(_ context.Context, name string) (iceio.File, error) { fs.mu.Lock() fs.opens[name]++ fs.mu.Unlock() - return fs.IO.Open(name) + return fs.IO.Open(context.Background(), name) } diff --git a/table/snapshot_manifest_cache_internal_test.go b/table/snapshot_manifest_cache_internal_test.go index 662646fe6..76f074205 100644 --- a/table/snapshot_manifest_cache_internal_test.go +++ b/table/snapshot_manifest_cache_internal_test.go @@ -264,7 +264,7 @@ type blockingSnapshotManifestIO struct { opens map[string]int } -func (fs *blockingSnapshotManifestIO) Open(name string) (iceio.File, error) { +func (fs *blockingSnapshotManifestIO) Open(_ context.Context, name string) (iceio.File, error) { if name == fs.blockedPath { fs.once.Do(func() { close(fs.started) }) <-fs.release @@ -274,7 +274,7 @@ func (fs *blockingSnapshotManifestIO) Open(name string) (iceio.File, error) { fs.opens[name]++ fs.mu.Unlock() - return fs.IO.Open(name) + return fs.IO.Open(context.Background(), name) } type notifyingContext struct { @@ -724,9 +724,9 @@ type contextBoundSnapshotManifestIO struct { blockedPath string } -func (fs *contextBoundSnapshotManifestIO) Open(name string) (iceio.File, error) { +func (fs *contextBoundSnapshotManifestIO) Open(_ context.Context, name string) (iceio.File, error) { if fs.blockedPath != "" && name != fs.blockedPath { - return fs.IO.Open(name) + return fs.IO.Open(context.Background(), name) } fs.startedOnce.Do(func() { close(fs.started) }) select { diff --git a/table/snapshot_producers.go b/table/snapshot_producers.go index 490fb2528..88057aa5f 100644 --- a/table/snapshot_producers.go +++ b/table/snapshot_producers.go @@ -292,7 +292,7 @@ func (of *overwriteFiles) rewriteManifest( if retErr == nil { return } - if cleanupErr := of.base.io.Remove(path); cleanupErr != nil { + if cleanupErr := of.base.io.Remove(context.Background(), path); cleanupErr != nil { retErr = errors.Join(retErr, fmt.Errorf("remove failed manifest %s: %w", path, cleanupErr)) } }() @@ -343,7 +343,7 @@ func (of *overwriteFiles) cleanupFilteredManifests(input []iceberg.ManifestFile, continue } seen[path] = struct{}{} - if err := of.base.io.Remove(path); err != nil { + if err := of.base.io.Remove(context.Background(), path); err != nil { cleanupErr = errors.Join(cleanupErr, fmt.Errorf("remove filtered manifest %s: %w", path, err)) } } @@ -424,7 +424,7 @@ func (of *overwriteFiles) filterAndRewriteParentManifests( func (sp *snapshotProducer) cleanupGeneratedManifests(manifests []iceberg.ManifestFile) error { var cleanupErr error for _, manifest := range manifests { - if err := sp.io.Remove(manifest.FilePath()); err != nil { + if err := sp.io.Remove(context.Background(), manifest.FilePath()); err != nil { cleanupErr = errors.Join(cleanupErr, fmt.Errorf("remove generated manifest %s: %w", manifest.FilePath(), err)) } } @@ -616,7 +616,7 @@ func (m *manifestMergeManager) removeOrphans(input, output []iceberg.ManifestFil if _, ok := inPaths[mf.FilePath()]; ok { continue } - if err := m.snap.io.Remove(mf.FilePath()); err != nil { + if err := m.snap.io.Remove(context.Background(), mf.FilePath()); err != nil { log.Printf("Warning: failed to delete orphaned merged manifest %s: %v", mf.FilePath(), err) } } @@ -849,7 +849,7 @@ func (sp *snapshotProducer) newManifestWriter(spec iceberg.PartitionSpec, opts . wr, err := iceberg.NewManifestWriter(sp.txn.meta.formatVersion, counter, spec, sp.txn.meta.CurrentSchema(), sp.snapshotID, opts...) if err != nil { - return nil, "", nil, nil, errors.Join(err, out.Close(), sp.io.Remove(path)) + return nil, "", nil, nil, errors.Join(err, out.Close(), sp.io.Remove(context.Background(), path)) } return wr, path, counter, out, nil @@ -1072,7 +1072,7 @@ func (sp *snapshotProducer) writeDeletedEntries(ctx context.Context, deleted []i if retErr == nil { return } - if cleanupErr := sp.io.Remove(path); cleanupErr != nil { + if cleanupErr := sp.io.Remove(ctx, path); cleanupErr != nil { retErr = errors.Join(retErr, fmt.Errorf("remove failed delete manifest %s: %w", path, cleanupErr)) } }() @@ -1489,7 +1489,7 @@ func writeManifestListFile( defer func() { if err != nil { - err = errors.Join(err, fs.Remove(path)) + err = errors.Join(err, fs.Remove(context.Background(), path)) } }() defer internal.CheckedClose(out, &err) diff --git a/table/snapshot_producers_bench_test.go b/table/snapshot_producers_bench_test.go index c02ea6092..ccb08088a 100644 --- a/table/snapshot_producers_bench_test.go +++ b/table/snapshot_producers_bench_test.go @@ -35,10 +35,10 @@ type delayedManifestIO struct { delay time.Duration } -func (io *delayedManifestIO) Open(name string) (iceio.File, error) { +func (io *delayedManifestIO) Open(_ context.Context, name string) (iceio.File, error) { time.Sleep(io.delay) - return io.memIO.Open(name) + return io.memIO.Open(context.Background(), name) } func (io *delayedManifestIO) Create(name string) (iceio.FileWriter, error) { diff --git a/table/snapshot_producers_test.go b/table/snapshot_producers_test.go index 47e60e13d..d19be7246 100644 --- a/table/snapshot_producers_test.go +++ b/table/snapshot_producers_test.go @@ -97,7 +97,7 @@ func newMemIO(limit int, err error) *memIO { } } -func (m *memIO) Open(name string) (iceio.File, error) { +func (m *memIO) Open(_ context.Context, name string) (iceio.File, error) { m.mu.Lock() data, ok := m.files[name] m.mu.Unlock() @@ -120,7 +120,7 @@ func (m *memIO) WriteFile(name string, content []byte) error { return nil } -func (m *memIO) Remove(name string) error { +func (m *memIO) Remove(_ context.Context, name string) error { m.mu.Lock() defer m.mu.Unlock() delete(m.files, name) @@ -452,7 +452,7 @@ func TestCommitV3RowLineageDeltaIncludesExistingRows(t *testing.T) { func readManifestListFromPath(t *testing.T, fs iceio.IO, path string) []iceberg.ManifestFile { t.Helper() - f, err := fs.Open(path) + f, err := fs.Open(context.Background(), path) require.NoError(t, err, "open manifest list: %s", path) defer f.Close() @@ -709,7 +709,7 @@ func newTrackingIO() *trackingIO { } } -func (t *trackingIO) Open(name string) (iceio.File, error) { +func (t *trackingIO) Open(_ context.Context, name string) (iceio.File, error) { data, ok := t.files[name] if !ok { return nil, fs.ErrNotExist @@ -734,7 +734,7 @@ func (t *trackingIO) WriteFile(name string, content []byte) error { return nil } -func (t *trackingIO) Remove(name string) error { +func (t *trackingIO) Remove(_ context.Context, name string) error { delete(t.files, name) return nil @@ -793,13 +793,13 @@ func (c *closeErrorIO) Create(name string) (iceio.FileWriter, error) { return writer, err } -func (c *closeErrorIO) Remove(name string) error { +func (c *closeErrorIO) Remove(_ context.Context, name string) error { c.removes = append(c.removes, name) if c.removeErr != nil { return c.removeErr } - return c.trackingIO.Remove(name) + return c.trackingIO.Remove(context.Background(), name) } func TestCommitManifestsCloseFailureReturnsNoUpdates(t *testing.T) { diff --git a/table/snapshots.go b/table/snapshots.go index 4d084332b..93df1b0c6 100644 --- a/table/snapshots.go +++ b/table/snapshots.go @@ -18,6 +18,7 @@ package table import ( + "context" "encoding/json" "errors" "fmt" @@ -406,7 +407,9 @@ func (s Snapshot) ValidateRowLineage() error { func (s Snapshot) Manifests(fio iceio.IO) (_ []iceberg.ManifestFile, err error) { if s.ManifestList != "" { - f, err := fio.Open(s.ManifestList) + // Manifests has no caller context; Open ignores ctx in every current IO + // backend, so Background is safe here. + f, err := fio.Open(context.Background(), s.ManifestList) if err != nil { return nil, fmt.Errorf("could not open manifest file: %w", err) } @@ -452,7 +455,7 @@ func embeddedManifestLength(fio iceio.IO, path string) (_ int64, err error) { return info.Size(), nil } - f, err := fio.Open(path) + f, err := fio.Open(context.Background(), path) if err != nil { return 0, fmt.Errorf("could not open embedded manifest %q: %w", path, err) } diff --git a/table/table.go b/table/table.go index 39ffa8649..e87030c85 100644 --- a/table/table.go +++ b/table/table.go @@ -625,7 +625,7 @@ func (t Table) doCommit(ctx context.Context, updates []Update, reqs []Requiremen return } for _, path := range orphanedManifests { - if removeErr := wfs.Remove(path); removeErr != nil { + if removeErr := wfs.Remove(ctx, path); removeErr != nil { log.Printf("Warning: failed to delete orphaned manifest list %s: %v", path, removeErr) } } @@ -1181,7 +1181,8 @@ func deleteOldMetadata(fs icebergio.IO, baseMeta, newMeta Metadata) { toRemove := internal.Difference(removedPrevious, currentMetadata) for _, file := range toRemove { - if err := fs.Remove(file); err != nil { + // deleteOldMetadata is best-effort cleanup with no caller context. + if err := fs.Remove(context.Background(), file); err != nil { // Log the error instead of raising it when deleting old metadata files, as an external entity like a compactor may have already deleted them log.Printf("Warning: Failed to delete old metadata file: %s error: %v", file, err) } @@ -1481,7 +1482,7 @@ func NewFromLocation( return nil, err } } else { - f, err := fsys.Open(metalocation) + f, err := fsys.Open(ctx, metalocation) if err != nil { return nil, err } diff --git a/table/table_test.go b/table/table_test.go index 7033e5bdd..3774a5320 100644 --- a/table/table_test.go +++ b/table/table_test.go @@ -2619,7 +2619,7 @@ func arrowTableWithNull() arrow.Table { } func (t *TableWritingTestSuite) validateManifestFileLength(fs iceio.IO, m iceberg.ManifestFile) { - f, err := fs.Open(m.FilePath()) + f, err := fs.Open(context.Background(), m.FilePath()) t.Require().NoError(err) defer f.Close() @@ -3416,7 +3416,7 @@ func (t *TableTestSuite) TestMetadataCompressionRoundTrip() { t.Require().NoError(err) // Read the metadata file and verify it's gzipped - file, err := fs.Open(metadataLoc) + file, err := fs.Open(context.Background(), metadataLoc) t.Require().NoError(err) defer file.Close() @@ -3475,7 +3475,7 @@ func (t *TableTestSuite) TestMetadataCompressionRoundTripZstd() { fs, err := tbl.FS(context.Background()) t.Require().NoError(err) - file, err := fs.Open(metadataLoc) + file, err := fs.Open(context.Background(), metadataLoc) t.Require().NoError(err) defer file.Close() diff --git a/table/transaction_test.go b/table/transaction_test.go index a2af4ac0b..98c3b765d 100644 --- a/table/transaction_test.go +++ b/table/transaction_test.go @@ -657,7 +657,7 @@ func (s *SparkIntegrationTestSuite) variantField(obj variant.ObjectValue, key st func (s *SparkIntegrationTestSuite) assertVariantFileShredded(tbl *table.Table, path string) { fs, err := tbl.FS(s.ctx) s.Require().NoError(err) - f, err := fs.Open(path) + f, err := fs.Open(s.ctx, path) s.Require().NoError(err) defer f.Close() diff --git a/table/updates.go b/table/updates.go index b9ad36b0c..2df115d19 100644 --- a/table/updates.go +++ b/table/updates.go @@ -803,7 +803,7 @@ func (u *removeSnapshotsUpdate) PostCommit(ctx context.Context, preTable *Table, var res error for _, f := range paths { - if err := prefs.Remove(f); err != nil { + if err := prefs.Remove(ctx, f); err != nil { res = errors.Join(res, err) } } diff --git a/table/updates_test.go b/table/updates_test.go index de500074c..c35a49780 100644 --- a/table/updates_test.go +++ b/table/updates_test.go @@ -52,20 +52,20 @@ func newTrackingCallsIO() *trackingCallsIO { } } -func (c *trackingCallsIO) Open(name string) (iceio.File, error) { +func (c *trackingCallsIO) Open(_ context.Context, name string) (iceio.File, error) { c.mu.Lock() c.openCount[name]++ c.mu.Unlock() - return c.trackingIO.Open(name) + return c.trackingIO.Open(context.Background(), name) } -func (c *trackingCallsIO) Remove(name string) error { +func (c *trackingCallsIO) Remove(_ context.Context, name string) error { c.mu.Lock() c.removeCount[name]++ c.mu.Unlock() - return c.trackingIO.Remove(name) + return c.trackingIO.Remove(context.Background(), name) } // writeManifest writes a v2 data manifest with a single ADDED diff --git a/table/write_records_test.go b/table/write_records_test.go index 6273c060b..cc2efb1ce 100644 --- a/table/write_records_test.go +++ b/table/write_records_test.go @@ -517,10 +517,10 @@ func (s *WriteRecordsTestSuite) TestTimestampNanosecondsShouldFloorNegativeValue type readOnlyFS struct{} -func (readOnlyFS) Open(name string) (iceio.File, error) { +func (readOnlyFS) Open(_ context.Context, name string) (iceio.File, error) { return nil, errors.New("not supported") } -func (readOnlyFS) Remove(name string) error { +func (readOnlyFS) Remove(_ context.Context, name string) error { return errors.New("not supported") } diff --git a/view/view.go b/view/view.go index a12be3938..23516f33b 100644 --- a/view/view.go +++ b/view/view.go @@ -84,7 +84,7 @@ func NewFromLocation( return nil, err } } else { - f, err := fsys.Open(metalocation) + f, err := fsys.Open(ctx, metalocation) if err != nil { return nil, err } diff --git a/view/view_test.go b/view/view_test.go index 977568757..6b389dbd0 100644 --- a/view/view_test.go +++ b/view/view_test.go @@ -110,7 +110,7 @@ func (t *ViewTestSuite) TestCreateViewJoinsTrailingSlashMetadataLocation() { fs, err := iceio.LoadFS(t.T().Context(), nil, metadataLocation) t.Require().NoError(err) - metadataFile, err := fs.Open(metadataLocation) + metadataFile, err := fs.Open(context.Background(), metadataLocation) t.Require().NoError(err) t.Require().NoError(metadataFile.Close()) }