SLoeuillet commented on code in PR #2079:
URL: https://github.com/apache/iceberg-go/pull/2079#discussion_r4163172475
##########
catalog/rest/vended_creds.go:
##########
@@ -171,33 +176,69 @@ 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 (blobfs does, so every gocloud backend), and a cached
IO
+ // built on a per-operation context would fail every later operation
with
+ // context.Canceled once that context is done. The refresher owns the
IO's
+ // lifetime instead, through ioCancel. fetchCreds above still honours
ctx.
+ //
+ // WithoutCancel keeps ctx's values, so the first caller's
request-scoped
+ // values stay attached to the shared IO.
+ ioCtx, ioCancel := context.WithCancel(context.WithoutCancel(ctx))
+
+ var (
+ newIO iceio.IO
+ expiresAt time.Time
+ issuedAt = v.issuedAt
+ )
+
if len(v.credentials) > 0 {
- prefixIO := newPrefixScopedIO(ctx, v.props, v.credentials)
+ prefixIO := newPrefixScopedIO(ioCtx, v.props, v.credentials)
prefixIO.nowFunc = v.now
- v.cachedIO = prefixIO
+ newIO = prefixIO
// Expiry is enforced by prefixScopedIO against only the
credential
// selected for each object location.
- v.expiresAt = time.Time{}
+ } else {
+ loaded, err := iceio.LoadFS(ioCtx, config, v.location)
+ if err != nil {
+ ioCancel()
- return v.cachedIO, nil
- }
+ if v.cachedIO == nil {
+ return nil, err
+ }
- newIO, err := iceio.LoadFS(ctx, config, v.location)
- if err != nil {
- if v.cachedIO == nil {
- return nil, err
+ return nil, fmt.Errorf("load filesystem with refreshed
credentials for %s: %w", v.location, err)
}
- return nil, fmt.Errorf("load filesystem with refreshed
credentials for %s: %w", v.location, err)
+ newIO = loaded
+ expiresAt = v.expiresAtFromConfig(config)
+ issuedAt = v.now()
}
- v.cachedIO = newIO
- v.expiresAt = v.expiresAtFromConfig(config)
- v.issuedAt = v.now()
+ v.replaceIO(newIO, ioCancel)
+ v.expiresAt = expiresAt
+ v.issuedAt = issuedAt
return v.cachedIO, nil
}
+// replaceIO installs io as the cached IO, then closes the one it supersedes
+// and cancels that one's context. An error closing the superseded IO is not
+// returned: the new IO is already in place, and the caller has no use for it.
+func (v *vendedCredentialRefresher) replaceIO(io iceio.IO, cancel
context.CancelFunc) {
+ oldIO, oldCancel := v.cachedIO, v.ioCancel
+ v.cachedIO, v.ioCancel = io, cancel
+
+ if oldIO != nil {
+ _ = closeOptionalIO(oldIO)
+ }
+
+ if oldCancel != nil {
+ oldCancel()
+ }
+}
Review Comment:
Done in d7ed212: `replaceIO` is swap-only, as you suggested. `close()` still
closes the current IO and cancels its context, and the `ioCancel` field comment
says so. Thanks for tracing the callers.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +931,168 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// contextCapturingScheme is a filesystem whose factory records every context
+// it is opened with, the way a filesystem that keeps its opening context
+// (blobfs) would, and counts how often each IO it handed out is closed.
+type contextCapturingScheme struct {
+ mu sync.Mutex
+ opened []context.Context
+ closes []*atomic.Int32
+}
+
+func registerContextCapturingScheme(t *testing.T, scheme string)
*contextCapturingScheme {
+ t.Helper()
+
+ c := &contextCapturingScheme{}
+
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ closes := &atomic.Int32{}
+
+ c.mu.Lock()
+ c.opened = append(c.opened, ctx)
+ c.closes = append(c.closes, closes)
+ c.mu.Unlock()
+
+ return &closeTrackingIO{IO: iceio.LocalFS{}, closeCount:
closes}, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return c
+}
+
+func (c *contextCapturingScheme) loads() int {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return len(c.opened)
+}
+
+func (c *contextCapturingScheme) ctx(i int) context.Context {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return c.opened[i]
+}
+
+func (c *contextCapturingScheme) closeCount(i int) int32 {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return c.closes[i].Load()
+}
+
+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. The
+// refresher owns the IO's lifetime instead: renewal closes the superseded IO
+// and ends its context.
+func TestVendedCredsCachedIODoesNotInheritCallerCancellation(t *testing.T) {
+ const scheme = "vended-ctx-detach-test"
+
+ fs := 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{},
"first"))
+ _, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.Equal(t, 1, fs.loads())
+ require.NoError(t, fs.ctx(0).Err(), "the cached IO must not be bound to
the caller's cancellation")
+ // WithoutCancel keeps values, so the first caller's values reach the
shared
+ // IO. This documents the current behaviour; it is a trade-off, not a
promise.
+ assert.Equal(t, "first", fs.ctx(0).Value(vendedCtxKey{}))
+
+ // Renewal rebuilds the IO on its own context and retires the old one.
+ r.expiresAt = time.Now().Add(-time.Minute)
+ ctx, cancel =
context.WithCancel(context.WithValue(context.Background(), vendedCtxKey{},
"second"))
+ _, err = r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.Equal(t, 2, fs.loads(), "renewal must rebuild the IO")
+ assert.NotSame(t, fs.ctx(0), fs.ctx(1), "the renewed IO gets a context
of its own")
+ require.NoError(t, fs.ctx(1).Err(), "the renewed IO must not be bound
to the caller's cancellation")
+ assert.Equal(t, "second", fs.ctx(1).Value(vendedCtxKey{}))
+
+ require.ErrorIs(t, fs.ctx(0).Err(), context.Canceled, "the superseded
IO's context is ended")
+ assert.Equal(t, int32(1), fs.closeCount(0), "the superseded IO is
closed")
Review Comment:
Inverted: after renewal the superseded IO's context is alive and it was not
closed. The test also keeps the IO from the first `loadFS`, renews, then writes
and reads through it. On 884d962 that read fails with `blob: Bucket has been
closed`.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +931,168 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// contextCapturingScheme is a filesystem whose factory records every context
+// it is opened with, the way a filesystem that keeps its opening context
+// (blobfs) would, and counts how often each IO it handed out is closed.
+type contextCapturingScheme struct {
+ mu sync.Mutex
+ opened []context.Context
+ closes []*atomic.Int32
+}
+
+func registerContextCapturingScheme(t *testing.T, scheme string)
*contextCapturingScheme {
+ t.Helper()
+
+ c := &contextCapturingScheme{}
+
+ iceio.Register(scheme, func(ctx context.Context, _ *url.URL, _
map[string]string) (iceio.IO, error) {
+ closes := &atomic.Int32{}
+
+ c.mu.Lock()
+ c.opened = append(c.opened, ctx)
+ c.closes = append(c.closes, closes)
+ c.mu.Unlock()
+
+ return &closeTrackingIO{IO: iceio.LocalFS{}, closeCount:
closes}, nil
+ })
+ t.Cleanup(func() { iceio.Unregister(scheme) })
+
+ return c
+}
+
+func (c *contextCapturingScheme) loads() int {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return len(c.opened)
+}
+
+func (c *contextCapturingScheme) ctx(i int) context.Context {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return c.opened[i]
+}
+
+func (c *contextCapturingScheme) closeCount(i int) int32 {
+ c.mu.Lock()
+ defer c.mu.Unlock()
+
+ return c.closes[i].Load()
+}
+
+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. The
+// refresher owns the IO's lifetime instead: renewal closes the superseded IO
+// and ends its context.
+func TestVendedCredsCachedIODoesNotInheritCallerCancellation(t *testing.T) {
+ const scheme = "vended-ctx-detach-test"
+
+ fs := 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{},
"first"))
+ _, err := r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.Equal(t, 1, fs.loads())
+ require.NoError(t, fs.ctx(0).Err(), "the cached IO must not be bound to
the caller's cancellation")
+ // WithoutCancel keeps values, so the first caller's values reach the
shared
+ // IO. This documents the current behaviour; it is a trade-off, not a
promise.
+ assert.Equal(t, "first", fs.ctx(0).Value(vendedCtxKey{}))
+
+ // Renewal rebuilds the IO on its own context and retires the old one.
+ r.expiresAt = time.Now().Add(-time.Minute)
+ ctx, cancel =
context.WithCancel(context.WithValue(context.Background(), vendedCtxKey{},
"second"))
+ _, err = r.loadFS(ctx)
+ require.NoError(t, err)
+ cancel()
+
+ require.Equal(t, 2, fs.loads(), "renewal must rebuild the IO")
+ assert.NotSame(t, fs.ctx(0), fs.ctx(1), "the renewed IO gets a context
of its own")
Review Comment:
Dropped. `loads() == 2` and the value check cover the rebuild.
##########
catalog/rest/vended_creds_test.go:
##########
@@ -928,3 +931,168 @@ func TestPrefixScopedIOPreservesReadContextCancellation(t
*testing.T) {
cancel()
require.ErrorIs(t, p.ctx.Err(), context.Canceled)
}
+
+// contextCapturingScheme is a filesystem whose factory records every context
+// it is opened with, the way a filesystem that keeps its opening context
+// (blobfs) would, and counts how often each IO it handed out is closed.
+type contextCapturingScheme struct {
Review Comment:
Done. The fake now returns a `blobfs.FileIO` over its own memblob bucket,
opened on the context it is given. The tests cancel the first caller's context
and then read through the cached IO, which fails on main with `context
canceled`.
--
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]