mirror of
https://github.com/k3s-io/kubernetes.git
synced 2026-07-23 02:00:12 +00:00
DRA kubelet: use TimedWorkersQueue
Conceptually TimedWorkersQueue is similar to the current code: it spawns goroutines and cancels them. Using it makes the code a bit shorter, even though the TimedWorkersQueue API could be a bit nicer and more consistent (key string vs. WorkerArgs as parameters). Depending on the tainteviction package is a bit odd, which is the reason why TimedWorkersQueue wasn't already used earlier. But there don't seem to be other implementations of this common problem. https://pkg.go.dev/k8s.io/client-go/util/workqueue#TypedDelayingInterface doesn't work because queue entries cannot be removed. This doesn't really solve the problem of tracking goroutines for wiping because TimedWorkersQueue doesn't support that. But not tracking is arguably better than doing it wrong and this only affects unit tests, so it should be okay.
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user