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]

Reply via email to