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]