@cryptotaxi247 / netdata-1 / commits / b7ff918d9

go.d.plugin: add docker service discovery (#17152)

Ilya Mashchenko committed Mar 14, 2024 at 11:27 UTC b7ff918d9161de68ef37e94e9845a81beba64b9a
16 files changed +852 -18
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/dockerd/docker.go new
+235
@@ -0,0 +1,235 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +package dockerd
4 +
5 +import (
6 + "context"
7 + "fmt"
8 + "log/slog"
9 + "net"
10 + "strconv"
11 + "strings"
12 + "time"
13 +
14 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
15 + "github.com/netdata/netdata/go/go.d.plugin/logger"
16 + "github.com/netdata/netdata/go/go.d.plugin/pkg/web"
17 +
18 + "github.com/docker/docker/api/types"
19 + typesContainer "github.com/docker/docker/api/types/container"
20 + docker "github.com/docker/docker/client"
21 + "github.com/ilyam8/hashstructure"
22 +)
23 +
24 +func NewDiscoverer(cfg Config) (*Discoverer, error) {
25 + tags, err := model.ParseTags(cfg.Tags)
26 + if err != nil {
27 + return nil, fmt.Errorf("parse tags: %v", err)
28 + }
29 +
30 + d := &Discoverer{
31 + Logger: logger.New().With(
32 + slog.String("component", "service discovery"),
33 + slog.String("discoverer", "docker"),
34 + ),
35 + cfgSource: cfg.Source,
36 + newDockerClient: func(addr string) (dockerClient, error) {
37 + return docker.NewClientWithOpts(docker.WithHost(addr))
38 + },
39 + addr: docker.DefaultDockerHost,
40 + listInterval: time.Second * 60,
41 + timeout: time.Second * 2,
42 + seenTggSources: make(map[string]bool),
43 + started: make(chan struct{}),
44 + }
45 +
46 + d.Tags().Merge(tags)
47 +
48 + if cfg.Timeout.Duration().Seconds() != 0 {
49 + d.timeout = cfg.Timeout.Duration()
50 + }
51 + if cfg.Address != "" {
52 + d.addr = cfg.Address
53 + }
54 +
55 + return d, nil
56 +}
57 +
58 +type Config struct {
59 + Source string
60 +
61 + Tags string `yaml:"tags"`
62 + Address string `yaml:"address"`
63 + Timeout web.Duration `yaml:"timeout"`
64 +}
65 +
66 +type (
67 + Discoverer struct {
68 + *logger.Logger
69 + model.Base
70 +
71 + dockerClient dockerClient
72 + newDockerClient func(addr string) (dockerClient, error)
73 + addr string
74 +
75 + cfgSource string
76 +
77 + listInterval time.Duration
78 + timeout time.Duration
79 + seenTggSources map[string]bool // [targetGroup.Source]
80 +
81 + started chan struct{}
82 + }
83 + dockerClient interface {
84 + NegotiateAPIVersion(context.Context)
85 + ContainerList(context.Context, typesContainer.ListOptions) ([]types.Container, error)
86 + Close() error
87 + }
88 +)
89 +
90 +func (d *Discoverer) String() string {
91 + return "sd:docker"
92 +}
93 +
94 +func (d *Discoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) {
95 + d.Info("instance is started")
96 + defer func() { d.cleanup(); d.Info("instance is stopped") }()
97 +
98 + close(d.started)
99 +
100 + if d.dockerClient == nil {
101 + client, err := d.newDockerClient(d.addr)
102 + if err != nil {
103 + d.Errorf("error on creating docker client: %v", err)
104 + return
105 + }
106 + d.dockerClient = client
107 + }
108 +
109 + d.dockerClient.NegotiateAPIVersion(ctx)
110 +
111 + if err := d.listContainers(ctx, in); err != nil {
112 + d.Error(err)
113 + return
114 + }
115 +
116 + tk := time.NewTicker(d.listInterval)
117 + defer tk.Stop()
118 +
119 + for {
120 + select {
121 + case <-ctx.Done():
122 + return
123 + case <-tk.C:
124 + if err := d.listContainers(ctx, in); err != nil {
125 + d.Warning(err)
126 + }
127 + }
128 + }
129 +}
130 +
131 +func (d *Discoverer) listContainers(ctx context.Context, in chan<- []model.TargetGroup) error {
132 + listCtx, cancel := context.WithTimeout(ctx, d.timeout)
133 + defer cancel()
134 +
135 + containers, err := d.dockerClient.ContainerList(listCtx, typesContainer.ListOptions{})
136 + if err != nil {
137 + return err
138 + }
139 +
140 + var tggs []model.TargetGroup
141 + seen := make(map[string]bool)
142 +
143 + for _, cntr := range containers {
144 + if tgg := d.buildTargetGroup(cntr); tgg != nil {
145 + tggs = append(tggs, tgg)
146 + seen[tgg.Source()] = true
147 + }
148 + }
149 +
150 + for src := range d.seenTggSources {
151 + if !seen[src] {
152 + tggs = append(tggs, &targetGroup{source: src})
153 + }
154 + }
155 + d.seenTggSources = seen
156 +
157 + select {
158 + case <-ctx.Done():
159 + case in <- tggs:
160 + }
161 +
162 + return nil
163 +}
164 +
165 +func (d *Discoverer) buildTargetGroup(cntr types.Container) model.TargetGroup {
166 + if len(cntr.Names) == 0 || cntr.NetworkSettings == nil || len(cntr.NetworkSettings.Networks) == 0 {
167 + return nil
168 + }
169 +
170 + tgg := &targetGroup{
171 + source: cntrSource(cntr),
172 + }
173 + if d.cfgSource != "" {
174 + tgg.source += fmt.Sprintf(",%s", d.cfgSource)
175 + }
176 +
177 + for netDriver, network := range cntr.NetworkSettings.Networks {
178 + // container with network mode host will be discovered by local-listeners
179 + for _, port := range cntr.Ports {
180 + tgt := &target{
181 + ID: cntr.ID,
182 + Name: strings.TrimPrefix(cntr.Names[0], "/"),
183 + Image: cntr.Image,
184 + Command: cntr.Command,
185 + Labels: mapAny(cntr.Labels),
186 + PrivatePort: strconv.Itoa(int(port.PrivatePort)),
187 + PublicPort: strconv.Itoa(int(port.PublicPort)),
188 + PublicPortIP: port.IP,
189 + PortProtocol: port.Type,
190 + NetworkMode: cntr.HostConfig.NetworkMode,
191 + NetworkDriver: netDriver,
192 + IPAddress: network.IPAddress,
193 + }
194 + tgt.Address = net.JoinHostPort(tgt.IPAddress, tgt.PrivatePort)
195 +
196 + hash, err := calcHash(tgt)
197 + if err != nil {
198 + continue
199 + }
200 +
201 + tgt.hash = hash
202 + tgt.Tags().Merge(d.Tags())
203 +
204 + tgg.targets = append(tgg.targets, tgt)
205 + }
206 + }
207 +
208 + return tgg
209 +}
210 +
211 +func (d *Discoverer) cleanup() {
212 + if d.dockerClient != nil {
213 + _ = d.dockerClient.Close()
214 + }
215 +}
216 +
217 +func cntrSource(cntr types.Container) string {
218 + name := strings.TrimPrefix(cntr.Names[0], "/")
219 + return fmt.Sprintf("discoverer=docker,container=%s,image=%s", name, cntr.Image)
220 +}
221 +
222 +func calcHash(obj any) (uint64, error) {
223 + return hashstructure.Hash(obj, nil)
224 +}
225 +
226 +func mapAny(src map[string]string) map[string]any {
227 + if src == nil {
228 + return nil
229 + }
230 + m := make(map[string]any, len(src))
231 + for k, v := range src {
232 + m[k] = v
233 + }
234 + return m
235 +}
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/dockerd/dockerd_test.go new
+161
@@ -0,0 +1,161 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +package dockerd
4 +
5 +import (
6 + "testing"
7 + "time"
8 +
9 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
10 +
11 + "github.com/docker/docker/api/types"
12 + typesNetwork "github.com/docker/docker/api/types/network"
13 +)
14 +
15 +func TestDiscoverer_Discover(t *testing.T) {
16 + tests := map[string]struct {
17 + createSim func() *discoverySim
18 + }{
19 + "add containers": {
20 + createSim: func() *discoverySim {
21 + nginx1 := prepareNginxContainer("nginx1")
22 + nginx2 := prepareNginxContainer("nginx2")
23 +
24 + sim := &discoverySim{
25 + dockerCli: func(cli dockerCli, _ time.Duration) {
26 + cli.addContainer(nginx1)
27 + cli.addContainer(nginx2)
28 + },
29 + wantGroups: []model.TargetGroup{
30 + &targetGroup{
31 + source: cntrSource(nginx1),
32 + targets: []model.Target{
33 + withHash(&target{
34 + ID: nginx1.ID,
35 + Name: nginx1.Names[0][1:],
36 + Image: nginx1.Image,
37 + Command: nginx1.Command,
38 + Labels: mapAny(nginx1.Labels),
39 + PrivatePort: "80",
40 + PublicPort: "8080",
41 + PublicPortIP: "0.0.0.0",
42 + PortProtocol: "tcp",
43 + NetworkMode: "default",
44 + NetworkDriver: "bridge",
45 + IPAddress: "192.0.2.0",
46 + Address: "192.0.2.0:80",
47 + }),
48 + },
49 + },
50 + &targetGroup{
51 + source: cntrSource(nginx2),
52 + targets: []model.Target{
53 + withHash(&target{
54 + ID: nginx2.ID,
55 + Name: nginx2.Names[0][1:],
56 + Image: nginx2.Image,
57 + Command: nginx2.Command,
58 + Labels: mapAny(nginx2.Labels),
59 + PrivatePort: "80",
60 + PublicPort: "8080",
61 + PublicPortIP: "0.0.0.0",
62 + PortProtocol: "tcp",
63 + NetworkMode: "default",
64 + NetworkDriver: "bridge",
65 + IPAddress: "192.0.2.0",
66 + Address: "192.0.2.0:80",
67 + }),
68 + },
69 + },
70 + },
71 + }
72 + return sim
73 + },
74 + },
75 + "remove containers": {
76 + createSim: func() *discoverySim {
77 + nginx1 := prepareNginxContainer("nginx1")
78 + nginx2 := prepareNginxContainer("nginx2")
79 +
80 + sim := &discoverySim{
81 + dockerCli: func(cli dockerCli, interval time.Duration) {
82 + cli.addContainer(nginx1)
83 + cli.addContainer(nginx2)
84 + time.Sleep(interval * 2)
85 + cli.removeContainer(nginx1.ID)
86 + },
87 + wantGroups: []model.TargetGroup{
88 + &targetGroup{
89 + source: cntrSource(nginx1),
90 + targets: nil,
91 + },
92 + &targetGroup{
93 + source: cntrSource(nginx2),
94 + targets: []model.Target{
95 + withHash(&target{
96 + ID: nginx2.ID,
97 + Name: nginx2.Names[0][1:],
98 + Image: nginx2.Image,
99 + Command: nginx2.Command,
100 + Labels: mapAny(nginx2.Labels),
101 + PrivatePort: "80",
102 + PublicPort: "8080",
103 + PublicPortIP: "0.0.0.0",
104 + PortProtocol: "tcp",
105 + NetworkMode: "default",
106 + NetworkDriver: "bridge",
107 + IPAddress: "192.0.2.0",
108 + Address: "192.0.2.0:80",
109 + }),
110 + },
111 + },
112 + },
113 + }
114 + return sim
115 + },
116 + },
117 + }
118 +
119 + for name, test := range tests {
120 + t.Run(name, func(t *testing.T) {
121 + sim := test.createSim()
122 + sim.run(t)
123 + })
124 + }
125 +}
126 +
127 +func prepareNginxContainer(name string) types.Container {
128 + return types.Container{
129 + ID: "id-" + name,
130 + Names: []string{"/" + name},
131 + Image: "nginx-image",
132 + ImageID: "nginx-image-id",
133 + Command: "nginx-command",
134 + Ports: []types.Port{
135 + {
136 + IP: "0.0.0.0",
137 + PrivatePort: 80,
138 + PublicPort: 8080,
139 + Type: "tcp",
140 + },
141 + },
142 + Labels: map[string]string{"key1": "value1"},
143 + HostConfig: struct {
144 + NetworkMode string `json:",omitempty"`
145 + }{
146 + NetworkMode: "default",
147 + },
148 + NetworkSettings: &types.SummaryNetworkSettings{
149 + Networks: map[string]*typesNetwork.EndpointSettings{
150 + "bridge": {IPAddress: "192.0.2.0"},
151 + },
152 + },
153 + }
154 +}
155 +
156 +func withHash(tgt *target) *target {
157 + tgt.hash, _ = calcHash(tgt)
158 + tags, _ := model.ParseTags("docker")
159 + tgt.Tags().Merge(tags)
160 + return tgt
161 +}
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/dockerd/sim_test.go new
+162
@@ -0,0 +1,162 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +package dockerd
4 +
5 +import (
6 + "context"
7 + "sort"
8 + "sync"
9 + "testing"
10 + "time"
11 +
12 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
13 +
14 + "github.com/docker/docker/api/types"
15 + typesContainer "github.com/docker/docker/api/types/container"
16 + "github.com/stretchr/testify/assert"
17 + "github.com/stretchr/testify/require"
18 +)
19 +
20 +type dockerCli interface {
21 + addContainer(cntr types.Container)
22 + removeContainer(id string)
23 +}
24 +
25 +type discoverySim struct {
26 + dockerCli func(cli dockerCli, interval time.Duration)
27 + wantGroups []model.TargetGroup
28 +}
29 +
30 +func (sim *discoverySim) run(t *testing.T) {
31 + d, err := NewDiscoverer(Config{
32 + Source: "",
33 + Tags: "docker",
34 + })
35 + require.NoError(t, err)
36 +
37 + mock := newMockDockerd()
38 +
39 + d.newDockerClient = func(addr string) (dockerClient, error) {
40 + return mock, nil
41 + }
42 + d.listInterval = time.Millisecond * 100
43 +
44 + seen := make(map[string]model.TargetGroup)
45 + ctx, cancel := context.WithCancel(context.Background())
46 + in := make(chan []model.TargetGroup)
47 + var wg sync.WaitGroup
48 +
49 + wg.Add(1)
50 + go func() {
51 + defer wg.Done()
52 + d.Discover(ctx, in)
53 + }()
54 +
55 + wg.Add(1)
56 + go func() {
57 + defer wg.Done()
58 + for {
59 + select {
60 + case <-ctx.Done():
61 + return
62 + case tggs := <-in:
63 + for _, tgg := range tggs {
64 + seen[tgg.Source()] = tgg
65 + }
66 + }
67 + }
68 + }()
69 +
70 + done := make(chan struct{})
71 + go func() {
72 + defer close(done)
73 + wg.Wait()
74 + }()
75 +
76 + select {
77 + case <-d.started:
78 + case <-time.After(time.Second * 3):
79 + require.Fail(t, "discovery failed to start")
80 + }
81 +
82 + sim.dockerCli(mock, d.listInterval)
83 + time.Sleep(time.Second)
84 +
85 + cancel()
86 +
87 + select {
88 + case <-done:
89 + case <-time.After(time.Second * 3):
90 + require.Fail(t, "discovery hasn't finished after cancel")
91 + }
92 +
93 + var tggs []model.TargetGroup
94 + for _, tgg := range seen {
95 + tggs = append(tggs, tgg)
96 + }
97 +
98 + sortTargetGroups(tggs)
99 + sortTargetGroups(sim.wantGroups)
100 +
101 + wantLen, gotLen := len(sim.wantGroups), len(tggs)
102 + assert.Equalf(t, wantLen, gotLen, "different len (want %d got %d)", wantLen, gotLen)
103 + assert.Equal(t, sim.wantGroups, tggs)
104 +
105 + assert.True(t, mock.negApiVerCalled, "NegotiateAPIVersion called")
106 + assert.True(t, mock.closeCalled, "Close called")
107 +}
108 +
109 +func newMockDockerd() *mockDockerd {
110 + return &mockDockerd{
111 + containers: make(map[string]types.Container),
112 + }
113 +}
114 +
115 +type mockDockerd struct {
116 + negApiVerCalled bool
117 + closeCalled bool
118 + mux sync.Mutex
119 + containers map[string]types.Container
120 +}
121 +
122 +func (m *mockDockerd) addContainer(cntr types.Container) {
123 + m.mux.Lock()
124 + defer m.mux.Unlock()
125 +
126 + m.containers[cntr.ID] = cntr
127 +}
128 +
129 +func (m *mockDockerd) removeContainer(id string) {
130 + m.mux.Lock()
131 + defer m.mux.Unlock()
132 +
133 + delete(m.containers, id)
134 +}
135 +
136 +func (m *mockDockerd) ContainerList(_ context.Context, _ typesContainer.ListOptions) ([]types.Container, error) {
137 + m.mux.Lock()
138 + defer m.mux.Unlock()
139 +
140 + var cntrs []types.Container
141 + for _, cntr := range m.containers {
142 + cntrs = append(cntrs, cntr)
143 + }
144 +
145 + return cntrs, nil
146 +}
147 +
148 +func (m *mockDockerd) NegotiateAPIVersion(_ context.Context) {
149 + m.negApiVerCalled = true
150 +}
151 +
152 +func (m *mockDockerd) Close() error {
153 + m.closeCalled = true
154 + return nil
155 +}
156 +
157 +func sortTargetGroups(tggs []model.TargetGroup) {
158 + if len(tggs) == 0 {
159 + return
160 + }
161 + sort.Slice(tggs, func(i, j int) bool { return tggs[i].Source() < tggs[j].Source() })
162 +}
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/dockerd/target.go new
+55
@@ -0,0 +1,55 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 +package dockerd
4 +
5 +import (
6 + "fmt"
7 +
8 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
9 +)
10 +
11 +type targetGroup struct {
12 + source string
13 + targets []model.Target
14 +}
15 +
16 +func (g *targetGroup) Provider() string { return "sd:docker" }
17 +func (g *targetGroup) Source() string { return g.source }
18 +func (g *targetGroup) Targets() []model.Target { return g.targets }
19 +
20 +type target struct {
21 + model.Base
22 +
23 + hash uint64
24 +
25 + ID string
26 + Name string
27 + Image string
28 + Command string
29 + Labels map[string]any
30 + PrivatePort string // Port on the container
31 + PublicPort string // Port exposed on the host
32 + PublicPortIP string // Host IP address that the container's port is mapped to
33 + PortProtocol string
34 + NetworkMode string
35 + NetworkDriver string
36 + IPAddress string
37 +
38 + Address string // "IPAddress:PrivatePort"
39 +}
40 +
41 +func (t *target) TUID() string {
42 + if t.PublicPort != "" {
43 + return fmt.Sprintf("%s_%s_%s_%s_%s_%s",
44 + t.Name, t.IPAddress, t.PublicPortIP, t.PortProtocol, t.PublicPort, t.PrivatePort)
45 + }
46 + if t.PrivatePort != "" {
47 + return fmt.Sprintf("%s_%s_%s_%s",
48 + t.Name, t.IPAddress, t.PortProtocol, t.PrivatePort)
49 + }
50 + return fmt.Sprintf("%s_%s", t.Name, t.IPAddress)
51 +}
52 +
53 +func (t *target) Hash() uint64 {
54 + return t.hash
55 +}
src/go/collectors/go.d.plugin/agent/discovery/sd/model/tags.go
+6
@@ -35,6 +35,12 @@ func (t Tags) Merge(tags Tags) {
35 }
36 }
37
38 +func (t Tags) Clone() Tags {
39 + ts := NewTags()
40 + ts.Merge(t)
41 + return ts
42 +}
43 +
44 func (t Tags) String() string {
45 ts := make([]string, 0, len(t))
46 for key := range t {
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/classify.go
+3 -1
@@ -46,10 +46,11 @@ type (
46 )
47
48 func (c *targetClassificator) classify(tgt model.Target) model.Tags {
49 + tgtTags := tgt.Tags().Clone()
50 var tags model.Tags
51
52 for i, rule := range c.rules {
52 - if !rule.sr.matches(tgt.Tags()) {
53 + if !rule.sr.matches(tgtTags) {
54 continue
55 }
56
@@ -70,6 +71,7 @@ func (c *targetClassificator) classify(tgt model.Target) model.Tags {
71
72 tags.Merge(rule.tags)
73 tags.Merge(match.tags)
74 + tgtTags.Merge(tags)
75 }
76 }
77
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/classify_test.go
+11 -2
@@ -14,14 +14,19 @@ import (
14
15 func TestTargetClassificator_classify(t *testing.T) {
16 config := `
17 -- selector: "rule1"
17 +- selector: "rule0"
18 + tags: "skip"
19 + match:
20 + - tags: "skip"
21 + expr: '{{ glob .Name "*" }}'
22 +- selector: "!skip rule1"
23 tags: "foo1"
24 match:
25 - tags: "bar1"
26 expr: '{{ glob .Name "mock*1*" }}'
27 - tags: "bar2"
28 expr: '{{ glob .Name "mock*2*" }}'
24 -- selector: "rule2"
29 +- selector: "!skip rule2"
30 tags: "foo2"
31 match:
32 - tags: "bar3"
@@ -56,6 +61,10 @@ func TestTargetClassificator_classify(t *testing.T) {
61 target: newMockTarget("mock123456", "rule1 rule2 rule3"),
62 wantTags: mustParseTags("foo1 foo2 foo3 bar1 bar2 bar3 bar4 bar5 bar6"),
63 },
64 + "applying labels after every rule": {
65 + target: newMockTarget("mock123456", "rule0 rule1 rule2 rule3"),
66 + wantTags: mustParseTags("skip foo3 bar5 bar6"),
67 + },
68 }
69
70 for name, test := range tests {
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/config.go
+3 -1
@@ -7,6 +7,7 @@ import (
7 "fmt"
8
9 "github.com/netdata/netdata/go/go.d.plugin/agent/confgroup"
10 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/dockerd"
11 "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/kubernetes"
12 "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/netlisteners"
13 )
@@ -24,6 +25,7 @@ type Config struct {
25 type DiscoveryConfig struct {
26 Discoverer string `yaml:"discoverer"`
27 NetListeners netlisteners.Config `yaml:"net_listeners"`
28 + Docker dockerd.Config `yaml:"docker"`
29 K8s []kubernetes.Config `yaml:"k8s"`
30 }
31
@@ -68,7 +70,7 @@ func validateDiscoveryConfig(config []DiscoveryConfig) error {
70 }
71 for _, cfg := range config {
72 switch cfg.Discoverer {
71 - case "net_listeners", "k8s":
73 + case "net_listeners", "docker", "k8s":
74 default:
75 return fmt.Errorf("unknown discoverer: '%s'", cfg.Discoverer)
76 }
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/pipeline.go
+15 -2
@@ -10,9 +10,11 @@ import (
10 "time"
11
12 "github.com/netdata/netdata/go/go.d.plugin/agent/confgroup"
13 + "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/dockerd"
14 "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/kubernetes"
15 "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/discoverer/netlisteners"
16 "github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
17 + "github.com/netdata/netdata/go/go.d.plugin/agent/hostinfo"
18 "github.com/netdata/netdata/go/go.d.plugin/logger"
19 )
20
@@ -78,14 +80,25 @@ func (p *Pipeline) registerDiscoverers(conf Config) error {
80 cfg.NetListeners.Source = conf.Source
81 td, err := netlisteners.NewDiscoverer(cfg.NetListeners)
82 if err != nil {
81 - return fmt.Errorf("failed to create 'net_listeners' discoverer: %v", err)
83 + return fmt.Errorf("failed to create '%s' discoverer: %v", cfg.Discoverer, err)
84 + }
85 + p.discoverers = append(p.discoverers, td)
86 + case "docker":
87 + if hostinfo.IsInsideK8sCluster() {
88 + p.Infof("not registering '%s' discoverer: disabled in k8s environment", cfg.Discoverer)
89 + continue
90 + }
91 + cfg.Docker.Source = conf.Source
92 + td, err := dockerd.NewDiscoverer(cfg.Docker)
93 + if err != nil {
94 + return fmt.Errorf("failed to create '%s' discoverer: %v", cfg.Discoverer, err)
95 }
96 p.discoverers = append(p.discoverers, td)
97 case "k8s":
98 for _, k8sCfg := range cfg.K8s {
99 td, err := kubernetes.NewKubeDiscoverer(k8sCfg)
100 if err != nil {
88 - return fmt.Errorf("failed to create 'k8s' discoverer: %v", err)
101 + return fmt.Errorf("failed to create '%s' discoverer: %v", cfg.Discoverer, err)
102 }
103 p.discoverers = append(p.discoverers, td)
104 }
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/pipeline_test.go
+2
@@ -25,6 +25,8 @@ func Test_defaultConfigs(t *testing.T) {
25 entries, err := os.ReadDir(dir)
26 require.NoError(t, err)
27
28 + require.NotEmpty(t, entries)
29 +
30 for _, e := range entries {
31 if strings.Contains(e.Name(), "prometheus") {
32 continue
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/promport.go
-2
@@ -11,8 +11,6 @@ var prometheusPortAllocations = map[int]string{
11 6060: "crowdsec",
12 7300: "midonet_agent",
13 8001: "netbox",
14 - 8080: "traefik",
15 - 8082: "trickster",
14 8088: "fawkes",
15 8089: "prom2teams",
16 8292: "phabricator_webhook_for_alertmanager",
src/go/collectors/go.d.plugin/agent/discovery/sd/pipeline/sim_test.go
+2 -2
@@ -33,10 +33,10 @@ func (sim discoverySim) run(t *testing.T) {
33 require.Nilf(t, err, "cfg unmarshal")
34
35 clr, err := newTargetClassificator(cfg.Classify)
36 - require.Nilf(t, err, "classify %v", err)
36 + require.Nil(t, err, "newTargetClassificator")
37
38 cmr, err := newConfigComposer(cfg.Compose)
39 - require.Nilf(t, err, "compose")
39 + require.Nil(t, err, "newConfigComposer")
40
41 mockClr := &mockClassificator{clr: clr}
42 mockCmr := &mockComposer{cmr: cmr}
src/go/collectors/go.d.plugin/agent/hostinfo/hostinfo.go
+10
@@ -5,6 +5,7 @@ package hostinfo
5 import (
6 "bytes"
7 "context"
8 + "os"
9 "os/exec"
10 "time"
11 )
@@ -27,3 +28,12 @@ func getHostname() string {
28
29 return string(bytes.TrimSpace(bs))
30 }
31 +
32 +var (
33 + envKubeHost = os.Getenv("KUBERNETES_SERVICE_HOST")
34 + envKubePort = os.Getenv("KUBERNETES_SERVICE_PORT")
35 +)
36 +
37 +func IsInsideK8sCluster() bool {
38 + return envKubeHost != "" && envKubePort != ""
39 +}
src/go/collectors/go.d.plugin/agent/setup.go
+2 -6
@@ -91,7 +91,7 @@ func (a *Agent) buildDiscoveryConf(enabled module.Registry) discovery.Config {
91 var readPaths, dummyPaths []string
92
93 if len(a.ModulesConfDir) == 0 {
94 - if isInsideK8sCluster() {
94 + if hostinfo.IsInsideK8sCluster() {
95 return discovery.Config{Registry: reg}
96 }
97 a.Info("modules conf dir not provided, will use default config for all enabled modules")
@@ -126,7 +126,7 @@ func (a *Agent) buildDiscoveryConf(enabled module.Registry) discovery.Config {
126 a.Debugf("looking for '%s' in %v", cfgName, a.ModulesConfDir)
127
128 path, err := a.ModulesConfDir.Find(cfgName)
129 - if isInsideK8sCluster() {
129 + if hostinfo.IsInsideK8sCluster() {
130 if err != nil {
131 a.Infof("not found '%s', won't use default (reading stock configs is disabled in k8s)", cfgName)
132 continue
@@ -196,13 +196,9 @@ func loadYAML(conf any, path string) error {
196 }
197
198 var (
199 - envKubeHost = os.Getenv("KUBERNETES_SERVICE_HOST")
200 - envKubePort = os.Getenv("KUBERNETES_SERVICE_PORT")
199 envNDStockConfigDir = os.Getenv("NETDATA_STOCK_CONFIG_DIR")
200 )
201
204 -func isInsideK8sCluster() bool { return envKubeHost != "" && envKubePort != "" }
205 -
202 func isStockConfig(path string) bool {
203 if envNDStockConfigDir == "" {
204 return false
src/go/collectors/go.d.plugin/config/go.d/sd/docker.conf new
+183
@@ -0,0 +1,183 @@
1 +name: 'docker'
2 +
3 +discover:
4 + - discoverer: docker
5 + docker:
6 + tags: "unknown"
7 + address: "unix:///var/run/docker.sock"
8 +
9 +classify:
10 + - name: "Skip"
11 + selector: "unknown"
12 + tags: "skip"
13 + match:
14 + - tags: "skip"
15 + expr: |
16 + {{ $netNOK := eq .NetworkMode "host" -}}
17 + {{ $protoNOK := not (eq .PortProtocol "tcp") -}}
18 + {{ $portNOK := empty .PrivatePort -}}
19 + {{ $addrNOK := or (empty .IPAddress) (glob .PublicPortIP "*:*") -}}
20 + {{ or $netNOK $protoNOK $portNOK $addrNOK }}
21 + - name: "Applications"
22 + selector: "!skip unknown"
23 + tags: "-unknown app"
24 + match:
25 + - tags: "apache"
26 + expr: '{{ match "sp" .Image "httpd httpd:* */apache */apache:* */apache2 */apache2:*" }}'
27 + - tags: "cockroachdb"
28 + expr: '{{ match "sp" .Image "cockroachdb/cockroach cockroachdb/cockroach:*" }}'
29 + - tags: "consul"
30 + expr: '{{ match "sp" .Image "consul consul:* */consul */consul:*" }}'
31 + - tags: "coredns"
32 + expr: '{{ match "sp" .Image "*/coredns */coredns:*" }}'
33 + - tags: "couchbase"
34 + expr: '{{ match "sp" .Image "couchbase couchbase:*" }}'
35 + - tags: "couchdb"
36 + expr: '{{ match "sp" .Image "couchdb couchdb:*" }}'
37 + - tags: "elasticsearch"
38 + expr: '{{ match "sp" .Image "elasticsearch elasticsearch:* */elasticsearch */elasticsearch:*" }}'
39 + - tags: "opensearch"
40 + expr: '{{ match "sp" .Image "*/opensearch */opensearch:*" }}'
41 + - tags: "lighttpd"
42 + expr: '{{ match "sp" .Image "*/lighttpd */lighttpd:*" }}'
43 + - tags: "mongodb"
44 + expr: '{{ match "sp" .Image "mongo mongo:* */mongodb */mongodb:*" }}'
45 + - tags: "mysql"
46 + expr: '{{ match "sp" .Image "mysql mysql:* */mysql */mysql:* mariadb mariadb:* */mariadb */mariadb:* percona percona:* */percona-mysql */percona-mysql:*" }}'
47 + - tags: "nginx"
48 + expr: '{{ match "sp" .Image "nginx nginx:*" }}'
49 + - tags: "pgbouncer"
50 + expr: '{{ match "sp" .Image "*/pgbouncer */pgbouncer:*" }}'
51 + - tags: "pika"
52 + expr: '{{ match "sp" .Image "pikadb/pika pikadb/pika:*" }}'
53 + - tags: "postgres"
54 + expr: '{{ match "sp" .Image "postgres postgres:* */postgres */postgres:* */postgresql */postgresql:*" }}'
55 + - tags: "proxysql"
56 + expr: '{{ match "sp" .Image "*/proxysql */proxysql:*" }}'
57 + - tags: "rabbitmq"
58 + expr: '{{ match "sp" .Image "rabbitmq rabbitmq:* */rabbitmq */rabbitmq:*" }}'
59 + - tags: "redis"
60 + expr: '{{ match "sp" .Image "redis redis:* */redis */redis:*" }}'
61 + - tags: "tengine"
62 + expr: '{{ match "sp" .Image "*/tengine */tengine:*" }}'
63 + - tags: "vernemq"
64 + expr: '{{ match "sp" .Image "*/vernemq */vernemq:*" }}'
65 + - tags: "zookeeper"
66 + expr: '{{ match "sp" .Image "*/zookeeper */zookeeper:*" }}'
67 +compose:
68 + - name: "Applications"
69 + selector: "app"
70 + config:
71 + - selector: "apache"
72 + template: |
73 + module: apache
74 + name: docker_{{.Name}}
75 + url: http://{{.Address}}/server-status?auto
76 + - selector: "cockroachdb"
77 + template: |
78 + module: cockroachdb
79 + name: docker_{{.Name}}
80 + url: http://{{.Address}}/_status/vars
81 + - selector: "consul"
82 + template: |
83 + module: consul
84 + name: docker_{{.Name}}
85 + url: http://{{.Address}}
86 + - selector: "coredns"
87 + template: |
88 + module: coredns
89 + name: docker_{{.Name}}
90 + url: http://{{.Address}}/metrics
91 + - selector: "coredns"
92 + template: |
93 + module: coredns
94 + name: docker_{{.Name}}
95 + url: http://{{.Address}}/metrics
96 + - selector: "couchbase"
97 + template: |
98 + module: couchbase
99 + name: docker_{{.Name}}
100 + url: http://{{.Address}}
101 + - selector: "couchdb"
102 + template: |
103 + module: couchdb
104 + name: docker_{{.Name}}
105 + url: http://{{.Address}}
106 + - selector: "elasticsearch"
107 + template: |
108 + module: elasticsearch
109 + name: docker_{{.Name}}
110 + url: http://{{.Address}}
111 + - selector: "opensearch"
112 + template: |
113 + module: elasticsearch
114 + name: docker_{{.Name}}
115 + url: https://{{.Address}}
116 + tls_skip_verify: yes
117 + username: admin
118 + password: admin
119 + - selector: "lighttpd"
120 + template: |
121 + module: lighttpd
122 + name: docker_{{.Name}}
123 + url: http://{{.Address}}/server-status?auto
124 + - selector: "mongodb"
125 + template: |
126 + module: mongodb
127 + name: docker_{{.Name}}
128 + uri: mongodb://{{.Address}}
129 + - selector: "mysql"
130 + template: |
131 + module: mysql
132 + name: docker_{{.Name}}
133 + dsn: netdata@tcp({{.IPAddress}}:{{.PrivatePort}})/
134 + - selector: "nginx"
135 + template: |
136 + module: nginx
137 + name: docker_{{.Name}}
138 + url: http://{{.Address}}/stub_status
139 + - selector: "pgbouncer"
140 + template: |
141 + module: pgbouncer
142 + name: docker_{{.Name}}
143 + dsn: postgres://netdata:postgres@{{.IPAddress}}:{{.PrivatePort}}/pgbouncer
144 + - selector: "pika"
145 + template: |
146 + module: pika
147 + name: docker_{{.Name}}
148 + address: redis://@{{.IPAddress}}:{{.PrivatePort}}
149 + - selector: "postgres"
150 + template: |
151 + module: postgres
152 + name: docker_{{.Name}}
153 + dsn: postgres://netdata:postgres@{{.IPAddress}}:{{.PrivatePort}}/postgres
154 + - selector: "proxysql"
155 + template: |
156 + module: proxysql
157 + name: docker_{{.Name}}
158 + dsn: stats:stats@tcp({{.IPAddress}}:{{.PrivatePort}})/
159 + - selector: "rabbitmq"
160 + template: |
161 + module: rabbitmq
162 + name: docker_{{.Name}}
163 + url: http://{{.Address}}
164 + - selector: "redis"
165 + template: |
166 + module: redis
167 + name: docker_{{.Name}}
168 + address: redis://@{{.IPAddress}}:{{.PrivatePort}}
169 + - selector: "tengine"
170 + template: |
171 + module: tengine
172 + name: docker_{{.Name}}
173 + url: http://{{.Address}}/us
174 + - selector: "vernemq"
175 + template: |
176 + module: vernemq
177 + name: docker_{{.Name}}
178 + url: http://{{.Address}}/metrics
179 + - selector: "zookeeper"
180 + template: |
181 + module: vernemq
182 + name: docker_{{.Name}}
183 + address: {{.Address}}
src/go/collectors/go.d.plugin/config/go.d/sd/net_listeners.conf
+2 -2
@@ -67,7 +67,7 @@ classify:
67 - tags: "mysql"
68 expr: '{{ and (eq .Port "3306") (eq .Comm "mysqld" "mariadb") }}'
69 - tags: "nginx"
70 - expr: '{{ and (eq .Port "80" "8080") (eq .Comm "nginx") }}'
70 + expr: '{{ and (eq .Port "80" "8080") (eq .Comm "nginx" "nginx:") }}'
71 - tags: "ntpd"
72 expr: '{{ and (eq .Port "123") (eq .Comm "ntpd") }}'
73 - tags: "openvpn"
@@ -280,7 +280,7 @@ compose:
280 template: |
281 module: nginx
282 name: local
283 - url: http://localhost:{{.Port}}/basic_status
283 + url: http://localhost:{{.Port}}/stub_status
284 - selector: "ntpd"
285 template: |
286 module: ntpd