diff --git a/pkg/kubelet/cm/dra/plugin/dra_plugin_manager.go b/pkg/kubelet/cm/dra/plugin/dra_plugin_manager.go index 50eedce9b1a..ae5cc972bc9 100644 --- a/pkg/kubelet/cm/dra/plugin/dra_plugin_manager.go +++ b/pkg/kubelet/cm/dra/plugin/dra_plugin_manager.go @@ -38,6 +38,7 @@ import ( "k8s.io/klog/v2" drapbv1alpha4 "k8s.io/kubelet/pkg/apis/dra/v1alpha4" drapbv1beta1 "k8s.io/kubelet/pkg/apis/dra/v1beta1" + timedworkers "k8s.io/kubernetes/pkg/controller/tainteviction" // TODO (?): move this common helper somewhere else? "k8s.io/kubernetes/pkg/kubelet/pluginmanager/cache" "k8s.io/utils/ptr" ) @@ -60,30 +61,20 @@ type DRAPluginManager struct { getNode func() (*v1.Node, error) wipingDelay time.Duration - // TODO: replace pendingWipes with some kind of workqueue. - // As it stands, WaitGroup suffers from a data race: - // - Queueing a new wiping creates a goroutine and adds to - // to wg. - // - Concurrently, wg.Wait reads from it. - // - // This is not allowed, all wg.Adds must come before wg.Wait. - // - // This race can be triggered with - // go test -count=10 -race ./... wg sync.WaitGroup mutex sync.RWMutex // driver name -> DRAPlugin in the order in which they got added store map[string][]*monitoredPlugin - // pendingWipes maps a driver name to a cancel function for - // wiping of that plugin's ResourceSlices. Entries get added - // in DeRegisterPlugin and check in RegisterPlugin. If - // wiping is pending during RegisterPlugin, it gets canceled. + // pendingWipes tracks at which time ResourceSlices for a + // DRA driver should be removed. The removal then happens in + // the background in a callback function that is invoked + // by the TimedWorkerQueue. // - // Must use pointers to functions because the entries have to - // be comparable. - pendingWipes map[string]*context.CancelCauseFunc + // TimedWorkerQueue uses namespace/name as key. We use + // the driver name as name with no namespace. + pendingWipes *timedworkers.TimedWorkerQueue } var _ cache.PluginHandler = &DRAPluginManager{} @@ -127,7 +118,10 @@ func (m *monitoredPlugin) HandleConn(_ context.Context, stats grpcstats.ConnStat default: return } - + if m.pm.backgroundCtx.Err() != nil { + // Shutting down, no longer interested in connection changes... + return + } logger := klog.FromContext(m.pm.backgroundCtx) m.pm.mutex.Lock() defer m.pm.mutex.Unlock() @@ -150,8 +144,11 @@ func NewDRAPluginManager(ctx context.Context, kubeClient kubernetes.Interface, g kubeClient: kubeClient, getNode: getNode, wipingDelay: wipingDelay, - pendingWipes: make(map[string]*context.CancelCauseFunc), } + pm.pendingWipes = timedworkers.CreateWorkerQueue(func(ctx context.Context, fireAt time.Time, args *timedworkers.WorkArgs) error { + pm.wipeResourceSlices(ctx, args.Object.Name) + return nil + }) // When kubelet starts up, no DRA driver has registered yet. None of // the drivers are usable until they come back, which might not happen @@ -167,51 +164,49 @@ func NewDRAPluginManager(ctx context.Context, kubeClient kubernetes.Interface, g ctx := pm.backgroundCtx logger := klog.LoggerWithName(klog.FromContext(ctx), "startup") ctx = klog.NewContext(ctx, logger) - pm.wipeResourceSlices(ctx, 0 /* no delay */, "" /* all drivers */) + pm.wipeResourceSlices(ctx, "" /* all drivers */) }() return pm } -// Stop cancels any remaining background activities and blocks until all goroutines have stopped. +// Stop cancels any remaining background activities and blocks until all goroutines have stopped, +// with one caveat: goroutines created dynamically for wiping ResourceSlices are not tracked. +// They won't do anything because of the context cancellation. func (pm *DRAPluginManager) Stop() { defer pm.wg.Wait() // Must run after unlocking our mutex. pm.mutex.Lock() defer pm.mutex.Unlock() + logger := klog.FromContext(pm.backgroundCtx) pm.cancel(errors.New("Stop was called")) // Close all connections, otherwise gRPC keeps doing things in the background. - for _, plugins := range pm.store { + // Also cancel all pending wiping. + for driverName, plugins := range pm.store { + workerArg := timedworkers.NewWorkArgs(driverName, "") + pm.pendingWipes.CancelWork(logger, workerArg.KeyFromWorkArgs()) for _, plugin := range plugins { if err := plugin.conn.Close(); err != nil { - klog.FromContext(pm.backgroundCtx).Error(err, "Closing gRPC connection", "driverName", plugin.driverName, "endpoint", plugin.endpoint) + logger.Error(err, "Closing gRPC connection", "driverName", plugin.driverName, "endpoint", plugin.endpoint) } } } } // wipeResourceSlices deletes ResourceSlices of the node, optionally just for a specific driver. -// Wiping will delay for a while and can be canceled by canceling the context. -func (pm *DRAPluginManager) wipeResourceSlices(ctx context.Context, delay time.Duration, driver string) { +// +// It gets called in a stand-alone goroutine at kubelet startup and as callback +// of a TimedWorkersQueue. In both cases the caller has no way of handling errors, +// so wipeResourceSlices must implement it's own retry mechanism. +// +// Can be canceled by canceling the context. +func (pm *DRAPluginManager) wipeResourceSlices(ctx context.Context, driver string) { if pm.kubeClient == nil { return } logger := klog.FromContext(ctx) - if delay != 0 { - // Before we start deleting, give the driver time to bounce back. - // Perhaps it got removed as part of a DaemonSet update and the - // replacement pod is about to start. - logger.V(4).Info("Starting to wait before wiping ResourceSlices", "delay", delay) - select { - case <-ctx.Done(): - logger.V(4).Info("Aborting wiping of ResourceSlices", "reason", context.Cause(ctx)) - case <-time.After(delay): - logger.V(4).Info("Starting to wipe ResourceSlices after waiting", "delay", delay) - } - } - backoff := wait.Backoff{ Duration: time.Second, Factor: 2, @@ -417,52 +412,39 @@ func (pm *DRAPluginManager) remove(driverName, endpoint string) { // sync must be called each time the information about a plugin changes. // The mutex must be locked for writing. func (pm *DRAPluginManager) sync(driverName string) { + if pm.kubeClient == nil { + // Cannot wipe. + return + } ctx := pm.backgroundCtx - logger := klog.FromContext(ctx) + logger := klog.FromContext(pm.backgroundCtx) + workerArgs := timedworkers.NewWorkArgs(driverName, "") // Is the DRA driver usable again? if pm.usable(driverName) { // Yes: cancel any pending ResourceSlice wiping for the DRA driver. - if cancel := pm.pendingWipes[driverName]; cancel != nil { - (*cancel)(errors.New("new plugin instance registered")) - delete(pm.pendingWipes, driverName) - } + pm.pendingWipes.CancelWork(logger, workerArgs.KeyFromWorkArgs()) return } - // No: prepare for canceling the background wiping. This needs to run - // in the context of the DRAPluginManager. + // No: ensure that we wipe ResourceSlices of the driver. + // If this was already queued earlier, the original timeout + // continues to apply because nothing changed. + if pm.pendingWipes.GetWorkerUnsafe(workerArgs.KeyFromWorkArgs()) != nil { + // Already queued or potentially already running. + // + // There's a small time-of-check-time-of-use race here, + // but that's fine: if wiping starts after we retrieve + // the pointer and before checking it, the work gets + // done, which is what we want. + return + } + now := time.Now() + fireAt := now.Add(pm.wipingDelay) logger = klog.LoggerWithName(logger, "driver-cleanup") logger = klog.LoggerWithValues(logger, "driverName", driverName) - ctx, cancel := context.WithCancelCause(pm.backgroundCtx) ctx = klog.NewContext(ctx, logger) - - // Clean up the ResourceSlices for the deleted DRAPlugin since it - // may have died without doing so itself and might never come - // back. - // - // May get canceled if the plugin comes back quickly enough. - if cancel := pm.pendingWipes[driverName]; cancel != nil { - (*cancel)(errors.New("plugin deregistered a second time")) - } - pm.pendingWipes[driverName] = &cancel - - pm.wg.Add(1) - go func() { - defer pm.wg.Done() - defer func() { - pm.mutex.Lock() - defer pm.mutex.Unlock() - - // Cancel our own context, but remove it from the map only if it - // is the current entry. Perhaps it already got replaced. - cancel(errors.New("wiping done")) - if pm.pendingWipes[driverName] == &cancel { - delete(pm.pendingWipes, driverName) - } - }() - pm.wipeResourceSlices(ctx, pm.wipingDelay, driverName) - }() + pm.pendingWipes.AddWork(ctx, timedworkers.NewWorkArgs(driverName, ""), now, fireAt) } // usable returns true if at least one endpoint is ready to handle gRPC calls for the DRA driver. diff --git a/pkg/kubelet/cm/dra/plugin/registration_test.go b/pkg/kubelet/cm/dra/plugin/registration_test.go index 6cd505df934..6c9e06dfd03 100644 --- a/pkg/kubelet/cm/dra/plugin/registration_test.go +++ b/pkg/kubelet/cm/dra/plugin/registration_test.go @@ -37,6 +37,7 @@ import ( cgotesting "k8s.io/client-go/testing" drapbv1alpha4 "k8s.io/kubelet/pkg/apis/dra/v1alpha4" drapb "k8s.io/kubelet/pkg/apis/dra/v1beta1" + timedworkers "k8s.io/kubernetes/pkg/controller/tainteviction" "k8s.io/kubernetes/test/utils/ktesting" ) @@ -333,6 +334,11 @@ func TestConnectionHandling(t *testing.T) { // Slice should get removed. requireNoSlices(tCtx) } else { + wipingIsPending := func() bool { + return draPlugins.pendingWipes.GetWorkerUnsafe(timedworkers.NewWorkArgs(driverName, "").KeyFromWorkArgs()) != nil + } + require.Eventuallyf(t, wipingIsPending, time.Minute, time.Second, "wiping should be queued for plugin %s", driverName) + // Start up gRPC server again. tCtx.Log("Restarting plugin gRPC server") teardown, err = setupFakeGRPCServer(service, endpoint) @@ -341,10 +347,7 @@ func TestConnectionHandling(t *testing.T) { // There shouldn't be any pending wipes for the plugin. require.Eventuallyf(t, func() bool { - draPlugins.mutex.Lock() - defer draPlugins.mutex.Unlock() - _, ok := draPlugins.pendingWipes[driverName] - return ok + return !wipingIsPending() }, time.Minute, time.Second, "wiping should be stopped for plugin %s", driverName) // Slice should still be there