laskoviymishka commented on code in PR #2079:
URL: https://github.com/apache/iceberg-go/pull/2079#discussion_r4148885187
##########
catalog/rest/vended_creds.go:
##########
@@ -171,8 +171,16 @@ func (v *vendedCredentialRefresher) loadFS(ctx
context.Context) (iceio.IO, error
maps.Copy(config, freshCreds)
}
+ // The IO is cached and shared by every later caller, so it must not
+ // inherit this caller's cancellation: a filesystem may keep the
context it
+ // is opened with (the GCS one does), so a cached IO built on a
+ // per-operation context would fail every later operation with
+ // context.Canceled once that context is done. Values are kept;
fetchCreds
+ // above still honours ctx.
+ ioCtx := context.WithoutCancel(ctx)
Review Comment:
Detaching cancellation here means `close()` is the only thing that can tear
down what the IO spawns. For a filesystem that keeps its opening context (the
GCS case this targets), background goroutines and credential refreshers now
live until `close()`, and renewal replaces the cached IO without closing the
old one (further down in `loadFS`), so the superseded one leaks for good, in
exactly the setup this PR fixes. Previously the caller's cancellation
eventually reaped that work; now nothing does.
I'd have the refresher own the cancel rather than fully detaching:
```go
ioCtx, ioCancel := context.WithCancel(context.WithoutCancel(ctx))
```
Store `ioCancel` on `v`, cancel the previous one and close the superseded IO
before replacing on renewal, and cancel it in `close()`. That keeps the fix and
closes the leak the detachment opens.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +928,87 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
Review Comment:
This test asserts `newPrefixScopedIO` keeps the caller's cancellation, but
the only production caller now hands it a detached context, so the contract it
names is dead for the vended path. That splits the detachment contract across
two spots and invites a later "fix" back to raw `ctx`. I'd either move the
`WithoutCancel` into `newPrefixScopedIO` so detachment lives in one place, or
add a line here noting the constructor is context-agnostic and `loadFS` owns
the detaching.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +928,87 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// registerContextCapturingScheme registers a filesystem whose factory records
+// the context it is opened with, the way a filesystem that keeps its opening
+// context (GCS) would.
+func registerContextCapturingScheme(t *testing.T, scheme string) func()
context.Context {
+ t.Helper()
+
+ var (
+ mu sync.Mutex
+ opened context.Context
+ )
+
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ mu.Lock()
+ opened = ctx
+ mu.Unlock()
+
+ return iceio.LocalFS{}, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return func() context.Context {
+ mu.Lock()
+ defer mu.Unlock()
+
+ return opened
+ }
+}
+
+type vendedCtxKey struct{}
+
+// The cached IO outlives the call that built it, so a caller that scopes one
+// operation to its own context (cancelled when the operation returns) must not
+// leave every later operation with a filesystem opened on a dead context.
+func TestVendedCredsCachedIODoesNotInheritCallerCancellation(t *testing.T) {
+ const scheme = "vended-ctx-detach-test"
+
+ opened := registerContextCapturingScheme(t, scheme)
+
+ r := newTestRefresher(func(context.Context, []string)
(iceberg.Properties, error) {
+ return iceberg.Properties{}, nil
+ })
+ r.location = scheme + "://bucket/tbl"
+
+ ctx, cancel :=
context.WithCancel(context.WithValue(context.Background(), vendedCtxKey{}, "v"))
+ _, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.NotNil(t, opened())
+ require.NoError(t, opened().Err(), "the cached IO must not be bound to
the caller's cancellation")
+ assert.Equal(t, "v", opened().Value(vendedCtxKey{}), "context values
still reach the filesystem")
Review Comment:
`WithoutCancel` keeps values, so the first caller's request-scoped values
(trace spans, loggers, per-request identity) stay attached to the shared IO and
reach every later, unrelated caller. Probably harmless for GCS, but asserting
it here as a guarantee bakes in a subtlety that's really a trade-off. I'd
soften this to document the behavior rather than present first-caller value
inheritance as intended.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +928,87 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// registerContextCapturingScheme registers a filesystem whose factory records
+// the context it is opened with, the way a filesystem that keeps its opening
+// context (GCS) would.
+func registerContextCapturingScheme(t *testing.T, scheme string) func()
context.Context {
+ t.Helper()
+
+ var (
+ mu sync.Mutex
+ opened context.Context
+ )
+
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ mu.Lock()
+ opened = ctx
+ mu.Unlock()
+
+ return iceio.LocalFS{}, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return func() context.Context {
+ mu.Lock()
+ defer mu.Unlock()
+
+ return opened
+ }
+}
+
+type vendedCtxKey struct{}
+
+// The cached IO outlives the call that built it, so a caller that scopes one
+// operation to its own context (cancelled when the operation returns) must not
+// leave every later operation with a filesystem opened on a dead context.
+func TestVendedCredsCachedIODoesNotInheritCallerCancellation(t *testing.T) {
+ const scheme = "vended-ctx-detach-test"
+
+ opened := registerContextCapturingScheme(t, scheme)
+
+ r := newTestRefresher(func(context.Context, []string)
(iceberg.Properties, error) {
+ return iceberg.Properties{}, nil
+ })
+ r.location = scheme + "://bucket/tbl"
+
+ ctx, cancel :=
context.WithCancel(context.WithValue(context.Background(), vendedCtxKey{}, "v"))
+ _, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.NotNil(t, opened())
+ require.NoError(t, opened().Err(), "the cached IO must not be bound to
the caller's cancellation")
+ assert.Equal(t, "v", opened().Value(vendedCtxKey{}), "context values
still reach the filesystem")
+
+ // A renewal rebuilds the IO; it must not inherit that caller's
cancellation either.
+ r.expiresAt = time.Now().Add(-time.Minute)
+ ctx, cancel = context.WithCancel(context.Background())
+ _, err = r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.NoError(t, opened().Err(), "the renewed IO must not be bound to
the caller's cancellation")
+}
+
+func TestVendedCredsPrefixScopedIODoesNotInheritCallerCancellation(t
*testing.T) {
+ const scheme = "vended-prefix-ctx-detach-test"
+
+ opened := registerContextCapturingScheme(t, scheme)
+
+ r := newTestRefresher(nil)
+ r.location = scheme + "://bucket/tbl"
+ r.credentials = []StorageCredential{{Prefix: scheme + "://bucket/",
Config: iceberg.Properties{}}}
+
+ ctx, cancel := context.WithCancel(context.Background())
+ fs, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ // prefixScopedIO opens its filesystems lazily, after the call that
built it returned.
+ _, err = fs.(*prefixScopedIO).filesystemFor(scheme +
"://bucket/tbl/data/file.parquet")
Review Comment:
Unchecked type assertion: if `loadFS` ever returns a different IO type this
panics and skips the cleanup output instead of failing cleanly. `p, ok :=
fs.(*prefixScopedIO); require.True(t, ok, ...)` reads better.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +928,87 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// registerContextCapturingScheme registers a filesystem whose factory records
+// the context it is opened with, the way a filesystem that keeps its opening
+// context (GCS) would.
+func registerContextCapturingScheme(t *testing.T, scheme string) func()
context.Context {
+ t.Helper()
+
+ var (
+ mu sync.Mutex
+ opened context.Context
+ )
+
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ mu.Lock()
+ opened = ctx
+ mu.Unlock()
+
+ return iceio.LocalFS{}, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return func() context.Context {
+ mu.Lock()
+ defer mu.Unlock()
+
+ return opened
+ }
+}
+
+type vendedCtxKey struct{}
+
+// The cached IO outlives the call that built it, so a caller that scopes one
+// operation to its own context (cancelled when the operation returns) must not
+// leave every later operation with a filesystem opened on a dead context.
+func TestVendedCredsCachedIODoesNotInheritCallerCancellation(t *testing.T) {
+ const scheme = "vended-ctx-detach-test"
+
+ opened := registerContextCapturingScheme(t, scheme)
+
+ r := newTestRefresher(func(context.Context, []string)
(iceberg.Properties, error) {
+ return iceberg.Properties{}, nil
+ })
+ r.location = scheme + "://bucket/tbl"
+
+ ctx, cancel :=
context.WithCancel(context.WithValue(context.Background(), vendedCtxKey{}, "v"))
+ _, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.NotNil(t, opened())
+ require.NoError(t, opened().Err(), "the cached IO must not be bound to
the caller's cancellation")
+ assert.Equal(t, "v", opened().Value(vendedCtxKey{}), "context values
still reach the filesystem")
+
+ // A renewal rebuilds the IO; it must not inherit that caller's
cancellation either.
+ r.expiresAt = time.Now().Add(-time.Minute)
+ ctx, cancel = context.WithCancel(context.Background())
+ _, err = r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.NoError(t, opened().Err(), "the renewed IO must not be bound to
the caller's cancellation")
Review Comment:
If renewal silently didn't rebuild (`needsRenewal` returned false),
`opened()` would still be the first, already-checked context and this passes
vacuously; nothing asserts the factory ran a second time. I'd record a call
count, or capture the first context before `cancel()` and assert the renewed
one is a distinct object. While here, the second ctx is built from
`Background()`, so the renewal path never checks values are kept; deriving it
with `WithValue` again would cover that too.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]