improvement(go.d.plugin): terminate on QUIT command (#19038)
Ilya Mashchenko committed
Nov 19, 2024 at 12:30 UTC
7b1ecbbb77859fbfbf7531f42e29bcb88e822b6b
4 files changed
+29
-13
src/go/plugin/go.d/agent/agent.go
+18
-7
@@ -62,6 +62,8 @@ type Agent struct {
62
Out io.Writer
63
64
api *netdataapi.API
65
+
66
+ quitCh chan struct{}
67
}
68
69
// New creates a new Agent.
@@ -83,6 +85,7 @@ func New(cfg Config) *Agent {
85
ModuleRegistry: module.DefaultRegistry,
86
Out: safewriter.Stdout,
87
api: netdataapi.New(safewriter.Stdout),
88
+ quitCh: make(chan struct{}),
89
}
90
}
91
@@ -100,17 +103,25 @@ func serve(a *Agent) {
103
var exit bool
104
105
for {
106
+ module.ObsoleteCharts(false)
107
+
108
ctx, cancel := context.WithCancel(context.Background())
109
110
wg.Add(1)
111
go func() { defer wg.Done(); a.run(ctx) }()
112
108
- switch sig := <-ch; sig {
109
- case syscall.SIGHUP:
110
- a.Infof("received %s signal (%d). Restarting running instance", sig, sig)
111
- default:
112
- a.Infof("received %s signal (%d). Terminating...", sig, sig)
113
- module.DontObsoleteCharts()
113
+ select {
114
+ case sig := <-ch:
115
+ switch sig {
116
+ case syscall.SIGHUP:
117
+ a.Infof("received %s signal (%d). Restarting running instance", sig, sig)
118
+ module.ObsoleteCharts(true)
119
+ default:
120
+ a.Infof("received %s signal (%d). Terminating...", sig, sig)
121
+ exit = true
122
+ }
123
+ case <-a.quitCh:
124
+ a.Infof("received QUIT command. Terminating...")
125
exit = true
126
}
127
@@ -209,7 +220,7 @@ func (a *Agent) run(ctx context.Context) {
220
var wg sync.WaitGroup
221
222
wg.Add(1)
212
- go func() { defer wg.Done(); fnMgr.Run(ctx) }()
223
+ go func() { defer wg.Done(); fnMgr.Run(ctx, a.quitCh) }()
224
225
wg.Add(1)
226
go func() { defer wg.Done(); jobMgr.Run(ctx, in) }()
src/go/plugin/go.d/agent/functions/manager.go
+8
-3
@@ -40,21 +40,21 @@ type Manager struct {
40
FunctionRegistry map[string]func(Function)
41
}
42
43
-func (m *Manager) Run(ctx context.Context) {
43
+func (m *Manager) Run(ctx context.Context, quitCh chan struct{}) {
44
m.Info("instance is started")
45
defer func() { m.Info("instance is stopped") }()
46
47
var wg sync.WaitGroup
48
49
wg.Add(1)
50
- go func() { defer wg.Done(); m.run(ctx) }()
50
+ go func() { defer wg.Done(); m.run(ctx, quitCh) }()
51
52
wg.Wait()
53
54
<-ctx.Done()
55
}
56
57
-func (m *Manager) run(ctx context.Context) {
57
+func (m *Manager) run(ctx context.Context, quitCh chan struct{}) {
58
for {
59
select {
60
case <-ctx.Done():
@@ -76,6 +76,11 @@ func (m *Manager) run(ctx context.Context) {
76
fn, err = parseFunctionWithPayload(ctx, line, m.input)
77
case line == "":
78
continue
79
+ case line == "QUIT":
80
+ if quitCh != nil {
81
+ quitCh <- struct{}{}
82
+ return
83
+ }
84
default:
85
m.Warningf("unexpected line: '%s'", line)
86
continue
src/go/plugin/go.d/agent/functions/manager_test.go
+1
-1
@@ -275,7 +275,7 @@ FUNCTION_PAYLOAD_END
275
276
done := make(chan struct{})
277
278
- go func() { defer close(done); mgr.Run(ctx) }()
278
+ go func() { defer close(done); mgr.Run(ctx, nil) }()
279
280
timeout := testTime + time.Second*2
281
tk := time.NewTimer(timeout)
src/go/plugin/go.d/agent/module/job.go
+2
-2
@@ -21,9 +21,9 @@ import (
21
var obsoleteLock = &sync.Mutex{}
22
var obsoleteCharts = true
23
24
-func DontObsoleteCharts() {
24
+func ObsoleteCharts(b bool) {
25
obsoleteLock.Lock()
26
- obsoleteCharts = false
26
+ obsoleteCharts = b
27
obsoleteLock.Unlock()
28
}
29