-
Notifications
You must be signed in to change notification settings - Fork 0
fix(credstore): support bounded keyring operations #79
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Open
robzolkos
wants to merge
5
commits into
main
Choose a base branch
from
fix-keyring-operation-timeouts
base: main
Could not load branches
Branch not found: {{ refName }}
Loading
Could not load tags
Nothing to show
Loading
Are you sure you want to change the base?
Some commits from the old base branch may be removed from the timeline,
and old review comments may become outdated.
Open
Changes from all commits
Commits
Show all changes
5 commits
Select commit
Hold shift + click to select a range
84a8b96
test: pin operation deadlines after a healthy keyring probe
robzolkos 60dca1a
Bound opted-in keyring operations without changing storage backends
robzolkos b3a9b86
no-mistakes(review): Clarify keyring timeout errors for reads and que…
robzolkos 86242a6
no-mistakes(document): Move credstore timeout contract into StoreOpti…
robzolkos afacf59
Prevent late keyring results from bypassing operation deadlines
robzolkos File filter
Filter by extension
Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
There are no files selected for viewing
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,73 @@ | ||
| package credstore | ||
|
|
||
| import ( | ||
| "context" | ||
| "fmt" | ||
| "time" | ||
| ) | ||
|
|
||
| // keyringOperation bounds both the wait behind another operation and the call | ||
| // itself. go-keyring has no cancellation API: after a timeout its worker may | ||
| // still finish. Refusing subsequent operations avoids overlapping an unknown | ||
| // write/delete and bounds abandoned workers to one per store. File storage is | ||
| // never substituted for a previously working keyring. | ||
| func (s *Store) keyringOperation(operation string, fn func() ([]byte, error)) ([]byte, error) { | ||
| if s.operationTimeout <= 0 { | ||
| return fn() | ||
| } | ||
| ctx, cancel := context.WithTimeout(context.Background(), s.operationTimeout) | ||
| defer cancel() | ||
| select { | ||
| case s.operationGate <- struct{}{}: | ||
| defer func() { <-s.operationGate }() | ||
| case <-ctx.Done(): | ||
| return nil, queueTimeoutError(operation, s.operationTimeout, ctx.Err()) | ||
| } | ||
| if s.operationErr != nil { | ||
| return nil, fmt.Errorf("keyring unavailable after an earlier timeout: %w", s.operationErr) | ||
| } | ||
| if err := operationDeadlineError(ctx); err != nil { | ||
| return nil, queueTimeoutError(operation, s.operationTimeout, err) | ||
| } | ||
|
|
||
| type result struct { | ||
| data []byte | ||
| err error | ||
| } | ||
| done := make(chan result, 1) | ||
| go func() { | ||
| data, err := fn() | ||
| done <- result{data, err} | ||
| }() | ||
| select { | ||
| case r := <-done: | ||
| return s.keyringOperationResult(ctx, operation, r.data, r.err) | ||
| case <-ctx.Done(): | ||
| return s.keyringOperationResult(ctx, operation, nil, ctx.Err()) | ||
| } | ||
| } | ||
|
|
||
| // keyringOperationResult is called with operationGate held, after dispatching | ||
| // the provider call. Deadline expiry wins even when the result was selected first. | ||
| func (s *Store) keyringOperationResult(ctx context.Context, operation string, data []byte, providerErr error) ([]byte, error) { | ||
| if err := operationDeadlineError(ctx); err != nil { | ||
| s.operationErr = fmt.Errorf("keyring %s timed out after %s (the operation may have completed or may still complete): %w", operation, s.operationTimeout, err) | ||
| return nil, s.operationErr | ||
| } | ||
| return data, providerErr | ||
| } | ||
|
|
||
| func operationDeadlineError(ctx context.Context) error { | ||
| if err := ctx.Err(); err != nil { | ||
| return err | ||
| } | ||
| // Timer delivery can lag behind the deadline when the scheduler is busy. | ||
| if deadline, ok := ctx.Deadline(); ok && !time.Now().Before(deadline) { | ||
| return context.DeadlineExceeded | ||
| } | ||
| return nil | ||
| } | ||
|
|
||
| func queueTimeoutError(operation string, timeout time.Duration, err error) error { | ||
| return fmt.Errorf("keyring %s timed out after %s waiting for another operation (not attempted): %w", operation, timeout, err) | ||
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,305 @@ | ||
| package credstore | ||
|
|
||
| import ( | ||
| "context" | ||
| "errors" | ||
| "os" | ||
| "sync" | ||
| "sync/atomic" | ||
| "testing" | ||
| "testing/synctest" | ||
| "time" | ||
|
|
||
| "github.com/stretchr/testify/assert" | ||
| "github.com/stretchr/testify/require" | ||
| ) | ||
|
|
||
| func swapKeyringOperations(t *testing.T, get func(string, string) (string, error), set func(string, string, string) error, del func(string, string) error) { | ||
| t.Helper() | ||
| oldGet, oldSet, oldDelete := keyringGet, keyringSet, keyringDelete | ||
| keyringGet, keyringSet, keyringDelete = get, set, del | ||
| t.Cleanup(func() { keyringGet, keyringSet, keyringDelete = oldGet, oldSet, oldDelete }) | ||
| } | ||
|
|
||
| // A responsive probe does not promise that the next real keyring operation | ||
| // will respond. Timing out must not serve a stale file, write plaintext, or | ||
| // start another operation alongside a write/delete whose outcome is unknown. | ||
| func TestKeyringOperationTimeout(t *testing.T) { | ||
| for _, operation := range []string{"load", "save", "delete", "migrate"} { | ||
| t.Run(operation, func(t *testing.T) { | ||
| dir := t.TempDir() | ||
| file := NewStore(StoreOptions{ForceFile: true, FallbackDir: dir}) | ||
| require.NoError(t, file.Save("work", []byte(`{"access_token":"stale-file-token"}`))) | ||
| before, err := os.ReadFile(file.credentialsPath()) | ||
| require.NoError(t, err) | ||
|
|
||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| var calls atomic.Int32 | ||
| release, finished := make(chan struct{}), make(chan struct{}) | ||
| block := func() { | ||
| calls.Add(1) | ||
| <-release | ||
| close(finished) | ||
| } | ||
| swapKeyringOperations(t, | ||
| func(string, string) (string, error) { block(); return "keyring-token", nil }, | ||
| func(string, string, string) error { block(); return nil }, | ||
| func(string, string) error { block(); return nil }, | ||
| ) | ||
| var releaseOnce sync.Once | ||
| t.Cleanup(func() { | ||
| releaseOnce.Do(func() { close(release) }) | ||
| if calls.Load() > 0 { | ||
| <-finished | ||
| } | ||
| }) | ||
| store := NewStore(StoreOptions{ServiceName: "test", FallbackDir: dir, OperationTimeout: 40 * time.Millisecond}) | ||
| result := make(chan error, 1) | ||
| go func() { | ||
| switch operation { | ||
| case "load": | ||
| data, loadErr := store.Load("work") | ||
| if len(data) != 0 { | ||
| loadErr = errors.New("timed-out read returned credential data") | ||
| } | ||
| result <- loadErr | ||
| case "save": | ||
| result <- store.Save("work", []byte("new-token")) | ||
| case "delete": | ||
| result <- store.Delete("work") | ||
| case "migrate": | ||
| result <- store.MigrateToKeyring() | ||
| } | ||
| }() | ||
| select { | ||
| case operationErr := <-result: | ||
| require.ErrorIs(t, operationErr, context.DeadlineExceeded) | ||
| assert.ErrorContains(t, operationErr, "may still complete") | ||
| assert.NotContains(t, operationErr.Error(), "not found", "a hung keyring is not a missing login") | ||
| case <-time.After(time.Second): | ||
| releaseOnce.Do(func() { close(release) }) | ||
| <-result | ||
| t.Fatal("keyring operation ignored OperationTimeout") | ||
| } | ||
|
|
||
| assert.True(t, store.UsingKeyring(), "a timed-out operation must not switch to plaintext") | ||
| assert.Empty(t, store.FallbackWarning()) | ||
| _, err = store.Load("work") | ||
| assert.ErrorIs(t, err, context.DeadlineExceeded) | ||
| assert.NotContains(t, err.Error(), "not found", "a hung keyring is not a missing login") | ||
| assert.ErrorIs(t, store.Save("work", []byte("newer-token")), context.DeadlineExceeded) | ||
| assert.ErrorIs(t, store.Delete("work"), context.DeadlineExceeded) | ||
| assert.EqualValues(t, 1, calls.Load(), "a stalled operation poisons this store, not the process") | ||
| after, err := os.ReadFile(file.credentialsPath()) | ||
| require.NoError(t, err) | ||
| assert.Equal(t, before, after, "timeout must neither change nor remove the fallback file") | ||
|
|
||
| // A late success is not permission to resume writes whose order | ||
| // relative to the timed-out write/delete cannot be established. | ||
| releaseOnce.Do(func() { close(release) }) | ||
| <-finished | ||
| assert.ErrorIs(t, store.Save("work", []byte("late-token")), context.DeadlineExceeded) | ||
| assert.EqualValues(t, 1, calls.Load()) | ||
| }) | ||
| } | ||
| } | ||
|
|
||
| func TestBoundedKeyringKeepsHealthyResultsAndErrors(t *testing.T) { | ||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| providerErr := errors.New("keyring locked") | ||
| swapKeyringOperations(t, | ||
| func(service, key string) (string, error) { | ||
| assert.Equal(t, "test", service) | ||
| assert.Equal(t, "test::work", key) | ||
| return "token", nil | ||
| }, | ||
| func(string, string, string) error { return providerErr }, | ||
| func(string, string) error { return nil }, | ||
| ) | ||
| store := NewStore(StoreOptions{ServiceName: "test", OperationTimeout: time.Second}) | ||
| data, err := store.Load("work") | ||
| require.NoError(t, err) | ||
| assert.Equal(t, "token", string(data)) | ||
| assert.ErrorIs(t, store.Save("work", data), providerErr) | ||
| assert.NoError(t, store.Delete("work"), "an ordinary provider error must not poison the store") | ||
| } | ||
|
|
||
| func TestWaitingForKeyringOperationAlsoHasADeadline(t *testing.T) { | ||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| var calls atomic.Int32 | ||
| swapKeyringOperations(t, | ||
| func(string, string) (string, error) { calls.Add(1); return "token", nil }, | ||
| func(string, string, string) error { return nil }, | ||
| func(string, string) error { return nil }, | ||
| ) | ||
| store := NewStore(StoreOptions{ServiceName: "test", OperationTimeout: 40 * time.Millisecond}) | ||
| store.operationGate <- struct{}{} | ||
| _, err := store.Load("work") | ||
| assert.ErrorIs(t, err, context.DeadlineExceeded) | ||
| assert.ErrorContains(t, err, "waiting for another operation") | ||
| assert.NotContains(t, err.Error(), "may still complete") | ||
| assert.Zero(t, calls.Load(), "a queued call must not reach the provider after its deadline") | ||
| <-store.operationGate | ||
| data, err := store.Load("work") | ||
| require.NoError(t, err, "timing out before calling the provider must not poison the store") | ||
| assert.Equal(t, "token", string(data)) | ||
| } | ||
|
|
||
| func TestZeroOperationTimeoutPreservesUnboundedCalls(t *testing.T) { | ||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| swapKeyringOperations(t, | ||
| func(string, string) (string, error) { return "token", nil }, | ||
| func(string, string, string) error { return nil }, | ||
| func(string, string) error { return nil }, | ||
| ) | ||
| store := NewStore(StoreOptions{ServiceName: "test"}) | ||
| data, err := store.Load("work") | ||
| require.NoError(t, err) | ||
| assert.Equal(t, "token", string(data)) | ||
| assert.NoError(t, store.Save("work", data)) | ||
| assert.NoError(t, store.Delete("work")) | ||
| } | ||
|
|
||
| func TestKeyringOperationResultRejectsExpiredDeadline(t *testing.T) { | ||
| for _, operation := range []string{"read", "write", "delete"} { | ||
| for _, timerDelivered := range []bool{true, false} { | ||
| name := operation + "/timer-pending" | ||
| if timerDelivered { | ||
| name = operation + "/timer-delivered" | ||
| } | ||
| t.Run(name, func(t *testing.T) { | ||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| var calls atomic.Int32 | ||
| swapKeyringOperations(t, nil, | ||
| func(string, string, string) error { calls.Add(1); return nil }, nil) | ||
| store := NewStore(StoreOptions{ServiceName: "test", OperationTimeout: time.Second}) | ||
|
|
||
| deadline := time.Now().Add(-time.Second) | ||
| var ctx context.Context | ||
| if timerDelivered { | ||
| var cancel context.CancelFunc | ||
| ctx, cancel = context.WithDeadline(context.Background(), deadline) | ||
| t.Cleanup(cancel) | ||
| require.ErrorIs(t, ctx.Err(), context.DeadlineExceeded) | ||
| } else { | ||
| parent, cancel := context.WithCancel(context.Background()) | ||
| t.Cleanup(cancel) | ||
| ctx = pendingDeadlineContext{Context: parent, deadline: deadline} | ||
| require.NoError(t, ctx.Err()) | ||
| } | ||
|
|
||
| // Model selecting the provider's successful result after expiry, | ||
| // whether or not the deadline timer has delivered its signal yet. | ||
| store.operationGate <- struct{}{} | ||
| data, err := store.keyringOperationResult(ctx, operation, []byte("late-token"), nil) | ||
| <-store.operationGate | ||
| require.ErrorIs(t, err, context.DeadlineExceeded) | ||
| assert.Empty(t, data) | ||
| assert.ErrorContains(t, err, "may have completed") | ||
| assert.ErrorIs(t, store.Save("work", []byte("retry")), context.DeadlineExceeded) | ||
| assert.Zero(t, calls.Load(), "a late result must not permit another provider call") | ||
| assert.True(t, store.UsingKeyring()) | ||
| }) | ||
| } | ||
| } | ||
| } | ||
|
|
||
| // pendingDeadlineContext models an expired deadline whose timer callback has | ||
| // not yet closed Done or set Err, as can happen when the scheduler is busy. | ||
| type pendingDeadlineContext struct { | ||
| context.Context | ||
| deadline time.Time | ||
| } | ||
|
|
||
| func (c pendingDeadlineContext) Deadline() (time.Time, bool) { | ||
| return c.deadline, true | ||
| } | ||
|
|
||
| func TestConcurrentKeyringCallsBehindStalledProvider(t *testing.T) { | ||
| synctest.Test(t, func(t *testing.T) { | ||
| stubProbe(t, func(string, time.Duration) error { return nil }) | ||
| file := NewStore(StoreOptions{ForceFile: true, FallbackDir: t.TempDir()}) | ||
| require.NoError(t, file.Save("work", []byte(`{"access_token":"stale-file-token"}`))) | ||
| before, err := os.ReadFile(file.credentialsPath()) | ||
| require.NoError(t, err) | ||
|
|
||
| started, release := make(chan struct{}), make(chan struct{}) | ||
| var startedOnce, releaseOnce sync.Once | ||
| var calls atomic.Int32 | ||
| var providers sync.WaitGroup | ||
| block := func() { | ||
| providers.Add(1) | ||
| defer providers.Done() | ||
| calls.Add(1) | ||
| startedOnce.Do(func() { close(started) }) | ||
| <-release | ||
| } | ||
| swapKeyringOperations(t, | ||
| func(string, string) (string, error) { block(); return "keyring-token", nil }, | ||
| func(string, string, string) error { block(); return nil }, | ||
| func(string, string) error { block(); return nil }, | ||
| ) | ||
| t.Cleanup(func() { | ||
| releaseOnce.Do(func() { close(release) }) | ||
| providers.Wait() | ||
| }) | ||
|
|
||
| const queued = 64 | ||
| const timeout = 40 * time.Millisecond | ||
| store := NewStore(StoreOptions{ServiceName: "test", FallbackDir: file.fallbackDir, OperationTimeout: timeout}) | ||
| type result struct { | ||
| data []byte | ||
| err error | ||
| } | ||
| results := make(chan result, queued+1) | ||
| go func() { results <- result{err: store.Save("work", []byte("new-token"))} }() | ||
| <-started | ||
| synctest.Wait() | ||
|
|
||
| for i := range queued { | ||
| go func() { | ||
| var r result | ||
| switch i % 3 { | ||
| case 0: | ||
| r.data, r.err = store.Load("work") | ||
| case 1: | ||
| r.err = store.Save("work", []byte("queued-token")) | ||
| case 2: | ||
| r.err = store.Delete("work") | ||
| } | ||
| results <- r | ||
| }() | ||
| } | ||
| // All queued callers are blocked behind the actual provider call | ||
| // before fake time advances to their deadlines. | ||
| synctest.Wait() | ||
| assert.Empty(t, results) | ||
| assert.EqualValues(t, 1, calls.Load()) | ||
| time.Sleep(timeout) | ||
| synctest.Wait() | ||
| require.Len(t, results, queued+1, "every caller must return while the provider is still blocked") | ||
| for range queued + 1 { | ||
| r := <-results | ||
| assert.ErrorIs(t, r.err, context.DeadlineExceeded) | ||
| assert.Empty(t, r.data, "a queued read must not return stale fallback credentials") | ||
| } | ||
| assert.EqualValues(t, 1, calls.Load()) | ||
| assert.True(t, store.UsingKeyring()) | ||
| assert.Empty(t, store.FallbackWarning()) | ||
|
|
||
| // A late successful write must not revive the store or permit any | ||
| // queued read/write/delete to reach the provider afterward. | ||
| releaseOnce.Do(func() { close(release) }) | ||
| providers.Wait() | ||
| synctest.Wait() | ||
| data, err := store.Load("work") | ||
| assert.ErrorIs(t, err, context.DeadlineExceeded) | ||
| assert.Empty(t, data) | ||
| assert.ErrorIs(t, store.Save("work", []byte("retry-token")), context.DeadlineExceeded) | ||
| assert.ErrorIs(t, store.Delete("work"), context.DeadlineExceeded) | ||
| assert.EqualValues(t, 1, calls.Load()) | ||
| after, err := os.ReadFile(file.credentialsPath()) | ||
| require.NoError(t, err) | ||
| assert.Equal(t, before, after, "the fallback file must remain untouched") | ||
| }) | ||
| } | ||
Oops, something went wrong.
Oops, something went wrong.
Add this suggestion to a batch that can be applied as a single commit.
This suggestion is invalid because no changes were made to the code.
Suggestions cannot be applied while the pull request is closed.
Suggestions cannot be applied while viewing a subset of changes.
Only one suggestion per line can be applied in a batch.
Add this suggestion to a batch that can be applied as a single commit.
Applying suggestions on deleted lines is not supported.
You must change the existing code in this line in order to create a valid suggestion.
Outdated suggestions cannot be applied.
This suggestion has been applied or marked resolved.
Suggestions cannot be applied from pending reviews.
Suggestions cannot be applied on multi-line comments.
Suggestions cannot be applied while the pull request is queued to merge.
Suggestion cannot be applied right now. Please check back later.
Uh oh!
There was an error while loading. Please reload this page.