master
go 131 lines 2.59 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package jobmgr
4
5 import (
6 "context"
7 "sync"
8
9 "github.com/netdata/netdata/go/plugins/plugin/framework/confgroup"
10 )
11
12 func newDiscoveredConfigsCache() *discoveredConfigs {
13 return &discoveredConfigs{
14 items: make(map[string]map[uint64]confgroup.Config),
15 }
16 }
17
18 func newRunningJobsCache() *runningJobs {
19 return &runningJobs{
20 mux: sync.Mutex{},
21 items: make(map[string]runtimeJob),
22 }
23 }
24
25 func newRetryingTasksCache() *retryingTasks {
26 return &retryingTasks{
27 items: make(map[string]*retryTask),
28 }
29 }
30
31 type (
32 discoveredConfigs struct {
33 // [Source][Hash]
34 items map[string]map[uint64]confgroup.Config
35 }
36
37 runningJobs struct {
38 mux sync.Mutex
39 // [cfg.FullName()]
40 items map[string]runtimeJob
41 }
42
43 retryingTasks struct {
44 // retryingTasks is intentionally lock-free. All access must remain serialized
45 // through manager-owned command flow and synchronous dyncfg handler callbacks.
46 // [cfg.UID()]
47 items map[string]*retryTask
48 }
49 retryTask struct {
50 cancel context.CancelFunc
51 }
52 )
53
54 func (c *discoveredConfigs) add(group *confgroup.Group) (added, removed []confgroup.Config) {
55 cfgs, ok := c.items[group.Source]
56 if !ok {
57 if len(group.Configs) == 0 {
58 return nil, nil
59 }
60 cfgs = make(map[uint64]confgroup.Config)
61 c.items[group.Source] = cfgs
62 }
63
64 seen := make(map[uint64]bool)
65
66 for _, cfg := range group.Configs {
67 hash := cfg.Hash()
68 seen[hash] = true
69
70 if _, ok := cfgs[hash]; ok {
71 continue
72 }
73
74 cfgs[hash] = cfg
75 added = append(added, cfg)
76 }
77
78 for hash, cfg := range cfgs {
79 if !seen[hash] {
80 delete(cfgs, hash)
81 removed = append(removed, cfg)
82 }
83 }
84
85 if len(cfgs) == 0 {
86 delete(c.items, group.Source)
87 }
88
89 return added, removed
90 }
91
92 func (c *runningJobs) lock() {
93 c.mux.Lock()
94 }
95 func (c *runningJobs) unlock() {
96 c.mux.Unlock()
97 }
98 func (c *runningJobs) add(fullName string, job runtimeJob) {
99 c.items[fullName] = job
100 }
101 func (c *runningJobs) remove(fullName string) {
102 delete(c.items, fullName)
103 }
104 func (c *runningJobs) lookup(fullName string) (runtimeJob, bool) {
105 j, ok := c.items[fullName]
106 return j, ok
107 }
108 func (c *runningJobs) snapshot() []runtimeJob {
109 c.mux.Lock()
110 defer c.mux.Unlock()
111
112 jobs := make([]runtimeJob, 0, len(c.items))
113 for _, job := range c.items {
114 jobs = append(jobs, job)
115 }
116 return jobs
117 }
118
119 func (c *retryingTasks) add(cfg confgroup.Config, retry *retryTask) {
120 c.items[cfg.UID()] = retry
121 }
122 func (c *retryingTasks) remove(cfg confgroup.Config) {
123 if v, ok := c.lookup(cfg); ok {
124 v.cancel()
125 }
126 delete(c.items, cfg.UID())
127 }
128 func (c *retryingTasks) lookup(cfg confgroup.Config) (*retryTask, bool) {
129 v, ok := c.items[cfg.UID()]
130 return v, ok
131 }