diff --git a/pkg/module_manager/module_manager.go b/pkg/module_manager/module_manager.go index 550ab9b1..c1e76094 100644 --- a/pkg/module_manager/module_manager.go +++ b/pkg/module_manager/module_manager.go @@ -1322,7 +1322,7 @@ func (mm *ModuleManager) UpdateModuleLastErrorAndNotify(module *modules.BasicMod // PushRunModuleTask pushes moduleRun task for a module into the main queue if there is no such a task for the module func (mm *ModuleManager) PushRunModuleTask(moduleName string, doModuleStartup bool) error { // check if there is already moduleRun task in the main queue for the module - if queueHasPendingModuleRunTaskWithStartup(mm.dependencies.TaskQueues.GetMain(), moduleName) { + if queueHasPendingModuleRunTask(mm.dependencies.TaskQueues.GetMain(), moduleName, doModuleStartup) { return nil } @@ -1682,9 +1682,15 @@ func (mm *ModuleManager) EnvironmentManagerEnabled() bool { return mm.environmentManager != nil } -// queueHasPendingModuleRunTaskWithStartup returns true if queue has pending tasks -// with the type "ModuleRun" related to the module "moduleName" and DoModuleStartup is set to true. -func queueHasPendingModuleRunTaskWithStartup(q *queue.TaskQueue, moduleName string) bool { +// queueHasPendingModuleRunTask returns true if the queue already holds a pending task with the type +// "ModuleRun" related to the module "moduleName" that covers a push made with doModuleStartup, so +// that pushing another one would only duplicate work. A pending task that runs the startup sequence +// covers one that does not; the reverse is not true, as the startup would then be dropped. +// +// Note that this must also hold for doModuleStartup=false on both sides, which is how every caller +// pushes today: answering only for pending tasks with startup lets identical tasks pile up, and the +// queue then reruns - and re-releases - the same module once per duplicate. +func queueHasPendingModuleRunTask(q *queue.TaskQueue, moduleName string, doModuleStartup bool) bool { if q == nil { return false } @@ -1692,7 +1698,7 @@ func queueHasPendingModuleRunTaskWithStartup(q *queue.TaskQueue, moduleName stri modules := modulesWithPendingTasks(q, task.ModuleRun) meta, has := modules[moduleName] - return has && meta.doStartup + return has && (meta.doStartup || !doModuleStartup) } func modulesWithPendingTasks(q *queue.TaskQueue, taskType sh_task.TaskType) map[string]struct{ doStartup bool } { @@ -1713,7 +1719,11 @@ func modulesWithPendingTasks(q *queue.TaskQueue, taskType sh_task.TaskType) map[ if t.GetType() == taskType { hm := task.HookMetadataAccessor(t) - modules[hm.ModuleName] = struct{ doStartup bool }{doStartup: hm.DoModuleStartup} + // a module may have several pending tasks: report the startup if any of them runs it, + // rather than letting the last task seen decide for all of them + meta := modules[hm.ModuleName] + meta.doStartup = meta.doStartup || hm.DoModuleStartup + modules[hm.ModuleName] = meta } }) diff --git a/pkg/module_manager/module_run_dedup_test.go b/pkg/module_manager/module_run_dedup_test.go new file mode 100644 index 00000000..e56cea0f --- /dev/null +++ b/pkg/module_manager/module_run_dedup_test.go @@ -0,0 +1,119 @@ +package module_manager + +import ( + "testing" + + "github.com/stretchr/testify/assert" + + "github.com/flant/addon-operator/pkg/task" + "github.com/flant/shell-operator/pkg/metric" + sh_task "github.com/flant/shell-operator/pkg/task" + "github.com/flant/shell-operator/pkg/task/queue" +) + +// Test_queueHasPendingModuleRunTask pins the guard that keeps PushRunModuleTask from stacking +// duplicate moduleRun tasks for one module. Every caller pushes with doModuleStartup=false, so a +// guard that only answers for pending tasks carrying the startup flag never fires for them: the +// main queue then fills with copies of the same task and the module is re-released once per copy. +func Test_queueHasPendingModuleRunTask(t *testing.T) { + const moduleName = "console" + + type pending struct { + module string + doStartup bool + } + + // The first task in the queue may already be running and is therefore not pending, so the queue + // is built with a leading task that is never the one under test. + newQueue := func(t *testing.T, tasks []pending) *queue.TaskQueue { + t.Helper() + + metricStorage := metric.NewStorageMock(t) + metricStorage.HistogramObserveMock.Set(func(_ string, _ float64, _ map[string]string, _ []float64) {}) + metricStorage.GaugeSetMock.Optional().Set(func(_ string, _ float64, _ map[string]string) {}) + + q := queue.NewTasksQueue("main", metricStorage) + q.AddLast(&sh_task.BaseTask{Type: task.ConvergeModules, Id: "running"}) + + for i, p := range tasks { + t := sh_task.NewTask(task.ModuleRun).WithMetadata(task.HookMetadata{ + ModuleName: p.module, + DoModuleStartup: p.doStartup, + }) + t.Id = string(rune('a' + i)) + q.AddLast(t) + } + + return q + } + + tests := []struct { + name string + pending []pending + doStartup bool + want bool + }{ + { + name: "nothing pending", + doStartup: false, + want: false, + }, + { + // the regression: identical pushes must collapse into one task + name: "pending task without startup covers a push without startup", + pending: []pending{{module: moduleName}}, + doStartup: false, + want: true, + }, + { + name: "pending task with startup covers a push without startup", + pending: []pending{{module: moduleName, doStartup: true}}, + doStartup: false, + want: true, + }, + { + name: "pending task with startup covers a push with startup", + pending: []pending{{module: moduleName, doStartup: true}}, + doStartup: true, + want: true, + }, + { + // the pending task does strictly less work, so dropping this push loses the startup + name: "pending task without startup does not cover a push with startup", + pending: []pending{{module: moduleName}}, + doStartup: true, + want: false, + }, + { + // the startup is already queued even though a later task does not carry it + name: "startup anywhere among the pending tasks counts", + pending: []pending{{module: moduleName, doStartup: true}, {module: moduleName}}, + doStartup: true, + want: true, + }, + { + name: "another module does not count", + pending: []pending{{module: "observability"}}, + doStartup: false, + want: false, + }, + { + // the only task for the module is the one at the head, which may already be running + name: "a task that is already running is not pending", + pending: nil, + doStartup: false, + want: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := queueHasPendingModuleRunTask(newQueue(t, tt.pending), moduleName, tt.doStartup) + assert.Equal(t, tt.want, got) + }) + } + + t.Run("nil queue", func(t *testing.T) { + assert.False(t, queueHasPendingModuleRunTask(nil, moduleName, false)) + }) +}