SLoeuillet commented on code in PR #2079:
URL: https://github.com/apache/iceberg-go/pull/2079#discussion_r4151909966
##########
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:
Done in 884d962. The refresher opens the IO on
`context.WithCancel(context.WithoutCancel(ctx))` and stores the cancel. Renewal
installs the new IO, then closes the superseded one and cancels its context. A
failed renewal cancels only the context it created and keeps the cached IO.
`close()` closes the IO and cancels its context. One trade-off: closing the
superseded IO also ends any operation still running on it. It is within the
renewal buffer of expiry by then, so I kept the close as you suggested.
##########
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:
Kept it in `loadFS`. With the refresher now owning a cancellable context,
`newPrefixScopedIO` has to keep the context it is given: that is how `close()`
reaches the filesystems it opens lazily. Moving `WithoutCancel` into the
constructor would cut that. I documented it on the constructor and on this test.
##########
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:
Softened: the test now documents first-caller value inheritance as current
behaviour and a trade-off, not a guarantee. The code comment says the same.
##########
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:
Good catch. The test now records every load: it asserts the factory ran
twice, the renewed context is a distinct object, it keeps the second caller's
values, and the superseded IO is closed and its context ended.
##########
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:
Fixed: `p, ok := loaded.(*prefixScopedIO); require.True(t, ok, ...)`.
--
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]