zeroshade commented on code in PR #2079:
URL: https://github.com/apache/iceberg-go/pull/2079#discussion_r4158471446


##########
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:
   Renewal closes the superseded IO and cancels its context while other callers 
can still hold it. `loadFS` hands the same IO to every caller, and the table 
code keeps it across later `fsF` calls. After gocloud `Bucket.Close()`, 
`NewRangeReader`/`NewWriter`/`Delete` return `blob: Bucket has been closed` 
(gocloud v0.46.0 `blob/blob.go:1034`, `:1187`, `:1284`). blobfs opens a range 
reader on every `ReadAt`, so files that are already open fail as well. The 
cancel also aborts reads and uploads that are in flight.
   
   In-tree callers this breaks:
   - `table/scanner.go:1181`: `manifestIOBatch.acquire` re-calls the factory at 
each batch boundary on purpose, as a renewal checkpoint. Sibling errgroup 
workers (`scanner.go:1280`, `:1435`) are still reading manifests through the 
previous IO, so a renewal there fails `PlanFiles`.
   - `table/rewrite_data_files.go:624`: `rewriteDataFilesPartial` holds `fs` 
across `ExecuteCompactionGroup` (`:647`, which reaches `fsF` via `ReadTasks` 
and `WriteRecords`). It then uses that `fs` for `CollectDeadPositionDeletes` 
(`:669`) and for cleanup (`:635`). A renewal mid-batch fails the batch, and 
because cleanup runs on the same closed IO, the new data files are orphaned.
   - `fsF` is shared by every table derived from the load 
(`table/table.go:824`, `table/transaction.go:3326`, `Scan.ioF`), so concurrent 
scans and writes, including lazy `ReadTasks` iterators, fail whenever another 
goroutine triggers renewal.
   
   "Within the renewal buffer" does not cover it. With no vended expiry, 
`expiresAtFromConfig` falls back to now+60m, so valid credentials get cut every 
~55 min where main never failed. With a known expiry, the buffer exists so 
in-flight work can finish on still-valid creds.
   
   The leak this guards against is not there for our backends: 
s3blob/gcsblob/azureblob `Close()` return nil (`Bucket.Close` is only a flag 
flip), `io/gocloud` starts no goroutines tied to the ctx, and 
`WithCancel(WithoutCancel(ctx))` registers with no parent, so an uncancelled 
one is just collected.
   
   The close also runs under `v.mu`, with the cancel after it. gocloud 
`Close()` takes the bucket write lock and waits for in-flight readers, so every 
other `loadFS` caller stalls behind them. 
`TestPrefixScopedIODoesNotHoldLockDuringFilesystemClose` exists to prevent that 
for `prefixScopedIO`.
   
   Fix: make this swap-only, as in f6423dc, and keep cancelling the current IO 
in `close()`. That is safe for plan-scoped IO because the scan's reader lease 
keeps it open while readers are active. Update the `ioCancel` field comment to 
match. If teardown for third-party IOs matters, retire the old IO after its 
credentials' hard expiry, outside the semaphore.
   
   ```suggestion
   // replaceIO installs io as the cached IO. The superseded IO is left open:
   // callers that loaded it earlier may still be using it, and its credentials
   // remain valid until their own expiry.
   func (v *vendedCredentialRefresher) replaceIO(io iceio.IO, cancel 
context.CancelFunc) {
        v.cachedIO, v.ioCancel = io, cancel
   }
   ```



##########
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:
   These two lines pin the teardown from the `replaceIO` comment. Invert them: 
after renewal, `fs.ctx(0).Err()` is nil and `closeCount(0)` is 0. Then add the 
property the table code depends on: keep the IO from the first `loadFS`, renew, 
and use it afterwards. With a memblob-backed `blobfs` scheme, that read would 
fail on this head 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:
   This only passes because `context.WithCancel` returns a `*cancelCtx`. 
Against a plain `WithoutCancel` context (f6423dc) it fails with `Both arguments 
must be pointers`; I ran it. It checks how the context type is implemented, not 
the renewal. `loads() == 2` plus the value check on :1023 already show a new IO 
was built on the second caller's context, so drop it.



##########
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:
   Optional: this fake records the context but never uses it, so the real root 
cause in #2078 is only covered by proxy: `blobfs` keeps the ctx and uses it in 
`Open`/`ReadAt`/`Create`. A scheme that returns a memblob-backed 
`blobfs.FileIO` built on the ctx it is given lets the test cancel the first 
caller's ctx, then open and read through the cached IO. Together with the 
renewal case on :1025, that covers both the bug and the regression.



-- 
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