fix(machine): improve context cancellation handling during watch and sync operations

This commit is contained in:
Pasha Sviderski
2026-06-24 13:41:14 +10:00
parent 83256b5bf7
commit 328d17599b
3 changed files with 9 additions and 5 deletions
+5 -1
View File
@@ -290,6 +290,7 @@ func (c *APIClient) resubscribeWithBackoffFn(id string) func(context.Context, ui
return nil return nil
} }
return func(ctx context.Context, fromChange uint64) (*Subscription, error) { return func(ctx context.Context, fromChange uint64) (*Subscription, error) {
boff := backoff.WithContext(c.newResubBackoff(), ctx)
return backoff.RetryWithData(func() (*Subscription, error) { return backoff.RetryWithData(func() (*Subscription, error) {
sub, err := c.ResubscribeContext(ctx, id, fromChange) sub, err := c.ResubscribeContext(ctx, id, fromChange)
if err != nil { if err != nil {
@@ -300,11 +301,14 @@ func (c *APIClient) resubscribeWithBackoffFn(id string) func(context.Context, ui
"id", id, "from_change", fromChange) "id", id, "from_change", fromChange)
return nil, backoff.Permanent(fmt.Errorf("resubscribe to %s: %w", id, err)) return nil, backoff.Permanent(fmt.Errorf("resubscribe to %s: %w", id, err))
} }
// Don't log retries triggered by context cancellation, the backoff will stop immediately.
if ctx.Err() == nil {
slog.Debug("Failed to resubscribe to Corrosion query. Retrying with backoff.", slog.Debug("Failed to resubscribe to Corrosion query. Retrying with backoff.",
"id", id, "from_change", fromChange, "err", err) "id", id, "from_change", fromChange, "err", err)
} }
}
return sub, err return sub, err
}, c.newResubBackoff()) }, boff)
} }
} }
+1 -1
View File
@@ -620,7 +620,7 @@ func (cc *clusterController) syncDockerContainers(ctx context.Context) error {
return nil return nil
} }
if err := backoff.Retry(watchAndSync, boff); err != nil { if err := backoff.Retry(watchAndSync, boff); err != nil {
if errors.Is(err, context.Canceled) { if ctx.Err() != nil {
return nil return nil
} }
return fmt.Errorf("watch and sync containers to cluster store: %w", err) return fmt.Errorf("watch and sync containers to cluster store: %w", err)
+1 -1
View File
@@ -139,7 +139,7 @@ func (c *Controller) WatchAndSyncContainers(ctx context.Context) error {
return fmt.Errorf("sync containers to cluster store: %w", err) return fmt.Errorf("sync containers to cluster store: %w", err)
} }
case err := <-errCh: case err := <-errCh:
if errors.Is(err, context.Canceled) { if ctx.Err() != nil {
return nil return nil
} }
return fmt.Errorf("receive Docker event: %w", err) return fmt.Errorf("receive Docker event: %w", err)