Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 16 additions & 6 deletions pkg/module_manager/module_manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down Expand Up @@ -1682,17 +1682,23 @@ 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
}

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 } {
Expand All @@ -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
}
})

Expand Down
119 changes: 119 additions & 0 deletions pkg/module_manager/module_run_dedup_test.go
Original file line number Diff line number Diff line change
@@ -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))
})
}
Loading