diff --git a/consumer_config.go b/consumer_config.go index 4beaf33..dec9e53 100644 --- a/consumer_config.go +++ b/consumer_config.go @@ -245,7 +245,7 @@ func (r *configRoulette) genConfigs(bestCfg *consumerConfig, queueEmpty bool) { r.oldBestCfg = bestCfg.Clone() r.addConfig(r.oldBestCfg) - if !hasFreeSystemResources() { + if !hasFreeSystemResources(r.opt.MinSystemResources) { internal.Logger.Println("taskq: system does not have enough free resources") return } diff --git a/go.mod b/go.mod index d41f34d..ab5e104 100644 --- a/go.mod +++ b/go.mod @@ -19,6 +19,7 @@ require ( github.com/onsi/gomega v1.10.3 github.com/satori/go.uuid v1.2.0 github.com/vmihailenco/msgpack/v5 v5.0.0 + github.com/stretchr/testify v1.6.1 golang.org/x/net v0.0.0-20201027133719-8eef5233e2a1 // indirect google.golang.org/protobuf v1.25.0 // indirect ) diff --git a/queue.go b/queue.go index d659860..2d647ce 100644 --- a/queue.go +++ b/queue.go @@ -59,6 +59,9 @@ type QueueOptions struct { // Optional message handler. The default is the global Tasks registry. Handler Handler + // Minimal system resources required to consider the consumer available to process the queue + MinSystemResources SystemResources + inited bool } @@ -118,6 +121,10 @@ func (opt *QueueOptions) Init() { if opt.Handler == nil { opt.Handler = &Tasks } + + if opt.MinSystemResources == (SystemResources{}) { + opt.MinSystemResources = NewDefaultSystemResources() + } } //------------------------------------------------------------------------------ diff --git a/sysinfo_linux.go b/sysinfo_linux.go deleted file mode 100644 index 1c0bc1a..0000000 --- a/sysinfo_linux.go +++ /dev/null @@ -1,31 +0,0 @@ -// +build linux - -package taskq - -import ( - "runtime" - - "github.com/capnm/sysinfo" -) - -func hasFreeSystemResources() bool { - si := sysinfo.Get() - free := si.FreeRam + si.BufferRam - - // at least 200MB of RAM is free - if free < 2e5 { - return false - } - - // at least 5% of RAM is free - if float64(free)/float64(si.TotalRam) < 0.05 { - return false - } - - // avg load is not too high - if si.Loads[0] > 1.5*float64(runtime.NumCPU()) { - return false - } - - return true -} diff --git a/sysinfo_other.go b/sysinfo_other.go deleted file mode 100644 index e6f7ed4..0000000 --- a/sysinfo_other.go +++ /dev/null @@ -1,7 +0,0 @@ -// +build !linux - -package taskq - -func hasFreeSystemResources() bool { - return true -} diff --git a/system.go b/system.go new file mode 100644 index 0000000..30691ca --- /dev/null +++ b/system.go @@ -0,0 +1,28 @@ +package taskq + +const ( + defaultSystemResourcesLoad1PerCPU float64 = 1.5 + defaultSystemResourcesMemoryFreeMB uint64 = 2e5 + defaultSystemResourcesMemoryFreePercentage uint64 = 5 +) + +// SystemResources represents system related values +type SystemResources struct { + // Maximum per CPU load at 1min intervals + Load1PerCPU float64 + + // Minimum free memory required in megabytes + MemoryFreeMB uint64 + + // Minimum free memory required in percentage + MemoryFreePercentage uint64 +} + +// NewDefaultSystemResources returns a new SystemResources struct with some default values +func NewDefaultSystemResources() SystemResources { + return SystemResources{ + Load1PerCPU: defaultSystemResourcesLoad1PerCPU, + MemoryFreeMB: defaultSystemResourcesMemoryFreeMB, + MemoryFreePercentage: defaultSystemResourcesMemoryFreePercentage, + } +} diff --git a/system_linux.go b/system_linux.go new file mode 100644 index 0000000..637f8ec --- /dev/null +++ b/system_linux.go @@ -0,0 +1,32 @@ +// +build linux + +package taskq + +import ( + "runtime" + + "github.com/capnm/sysinfo" + "github.com/vmihailenco/taskq/v3/internal" +) + +func hasFreeSystemResources(sr SystemResources) bool { + si := sysinfo.Get() + free := si.FreeRam + si.BufferRam + + if sr.Load1PerCPU > 0 && si.Loads[0] > sr.Load1PerCPU*float64(runtime.NumCPU()) { + internal.Logger.Println("taskq: consumer memory is lower than required") + return false + } + + if sr.MemoryFreeMB > 0 && free < sr.MemoryFreeMB { + internal.Logger.Println("taskq: consumer memory is lower than required") + return false + } + + if sr.MemoryFreePercentage > 0 && free/si.TotalRam < sr.MemoryFreePercentage/100 { + internal.Logger.Println("taskq: consumer memory is lower than required") + return false + } + + return true +} diff --git a/system_linux_test.go b/system_linux_test.go new file mode 100644 index 0000000..76456ba --- /dev/null +++ b/system_linux_test.go @@ -0,0 +1,12 @@ +package taskq + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestHasFreeSystemResources(t *testing.T) { + // TODO: Manage to mock capnm/sysinfo + assert.True(t, hasFreeSystemResources(SystemResources{})) +} diff --git a/system_other.go b/system_other.go new file mode 100644 index 0000000..588aaab --- /dev/null +++ b/system_other.go @@ -0,0 +1,7 @@ +// +build !linux + +package taskq + +func hasFreeSystemResources(_ SystemResources) bool { + return true +} diff --git a/system_other_test.go b/system_other_test.go new file mode 100644 index 0000000..2ca53e5 --- /dev/null +++ b/system_other_test.go @@ -0,0 +1,11 @@ +package taskq + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestHasFreeSystemResources(t *testing.T) { + assert.True(t, hasFreeSystemResources(SystemResources{})) +} diff --git a/system_test.go b/system_test.go new file mode 100644 index 0000000..8095bbf --- /dev/null +++ b/system_test.go @@ -0,0 +1,16 @@ +package taskq + +import ( + "testing" + + "github.com/stretchr/testify/assert" +) + +func TestNewDefaultSystemResources(t *testing.T) { + expectedValue := SystemResources{ + Load1PerCPU: 1.5, + MemoryFreeMB: 2e5, + MemoryFreePercentage: 5, + } + assert.Equal(t, expectedValue, NewDefaultSystemResources()) +}