master
go 215 lines 4.97 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package dockersd
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/plugins/logger"
15 "github.com/netdata/netdata/go/plugins/pkg/confopt"
16 "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/model"
17 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/dockerhost"
18
19 typesContainer "github.com/docker/docker/api/types/container"
20 docker "github.com/docker/docker/client"
21 )
22
23 func NewDiscoverer(cfg Config) (*Discoverer, error) {
24 d := &Discoverer{
25 Logger: logger.New().With(
26 slog.String("component", "service discovery"),
27 slog.String("discoverer", "docker"),
28 ),
29 cfgSource: cfg.Source,
30 newDockerClient: func(addr string) (dockerClient, error) {
31 return docker.NewClientWithOpts(docker.WithHost(addr))
32 },
33 addr: docker.DefaultDockerHost,
34 listInterval: time.Second * 60,
35 timeout: time.Second * 2,
36 seenTggSources: make(map[string]bool),
37 started: make(chan struct{}),
38 }
39
40 if addr := dockerhost.FromEnv(); addr != "" && d.addr == docker.DefaultDockerHost {
41 d.Infof("using docker host from environment: %s ", addr)
42 d.addr = addr
43 }
44
45 if cfg.Timeout.Duration() > 0 {
46 d.timeout = cfg.Timeout.Duration()
47 }
48 if cfg.Address != "" {
49 d.addr = cfg.Address
50 }
51
52 return d, nil
53 }
54
55 type Config struct {
56 Source string `yaml:"-" json:"-"`
57
58 Address string `yaml:"address,omitempty" json:"address,omitempty"`
59 Timeout confopt.Duration `yaml:"timeout,omitempty" json:"timeout,omitempty"`
60 }
61
62 type (
63 Discoverer struct {
64 *logger.Logger
65 model.Base
66
67 dockerClient dockerClient
68 newDockerClient func(addr string) (dockerClient, error)
69 addr string
70
71 cfgSource string
72
73 listInterval time.Duration
74 timeout time.Duration
75 seenTggSources map[string]bool // [targetGroup.Source]
76
77 started chan struct{}
78 }
79 dockerClient interface {
80 NegotiateAPIVersion(context.Context)
81 ContainerList(context.Context, typesContainer.ListOptions) ([]typesContainer.Summary, error)
82 Close() error
83 }
84 )
85
86 func (d *Discoverer) String() string {
87 return "sd:docker"
88 }
89
90 func (d *Discoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) {
91 d.Info("instance is started")
92 defer func() { d.cleanup(); d.Info("instance is stopped") }()
93
94 close(d.started)
95
96 if d.dockerClient == nil {
97 client, err := d.newDockerClient(d.addr)
98 if err != nil {
99 d.Errorf("error on creating docker client: %v", err)
100 return
101 }
102 d.dockerClient = client
103 }
104
105 d.dockerClient.NegotiateAPIVersion(ctx)
106
107 if err := d.listContainers(ctx, in); err != nil {
108 d.Error(err)
109 return
110 }
111
112 tk := time.NewTicker(d.listInterval)
113 defer tk.Stop()
114
115 for {
116 select {
117 case <-ctx.Done():
118 return
119 case <-tk.C:
120 if err := d.listContainers(ctx, in); err != nil {
121 d.Warning(err)
122 }
123 }
124 }
125 }
126
127 func (d *Discoverer) listContainers(ctx context.Context, in chan<- []model.TargetGroup) error {
128 listCtx, cancel := context.WithTimeout(ctx, d.timeout)
129 defer cancel()
130
131 containers, err := d.dockerClient.ContainerList(listCtx, typesContainer.ListOptions{})
132 if err != nil {
133 return err
134 }
135
136 var tggs []model.TargetGroup
137 seen := make(map[string]bool)
138
139 for _, cntr := range containers {
140 if tgg := d.buildTargetGroup(cntr); tgg != nil {
141 tggs = append(tggs, tgg)
142 seen[tgg.Source()] = true
143 }
144 }
145
146 for src := range d.seenTggSources {
147 if !seen[src] {
148 tggs = append(tggs, &targetGroup{source: src})
149 }
150 }
151 d.seenTggSources = seen
152
153 select {
154 case <-ctx.Done():
155 case in <- tggs:
156 }
157
158 return nil
159 }
160
161 func (d *Discoverer) buildTargetGroup(cntr typesContainer.Summary) model.TargetGroup {
162 if len(cntr.Names) == 0 || cntr.NetworkSettings == nil || len(cntr.NetworkSettings.Networks) == 0 {
163 return nil
164 }
165
166 tgg := &targetGroup{
167 source: cntrSource(cntr),
168 }
169 if d.cfgSource != "" {
170 tgg.source += fmt.Sprintf(",%s", d.cfgSource)
171 }
172
173 for netDriver, network := range cntr.NetworkSettings.Networks {
174 // container with network mode host will be discovered by local-listeners
175 for _, port := range cntr.Ports {
176 tgt := &target{
177 ID: cntr.ID,
178 Name: strings.TrimPrefix(cntr.Names[0], "/"),
179 Image: cntr.Image,
180 Command: cntr.Command,
181 Labels: model.MapAny(cntr.Labels),
182 PrivatePort: strconv.Itoa(int(port.PrivatePort)),
183 PublicPort: strconv.Itoa(int(port.PublicPort)),
184 PublicPortIP: port.IP,
185 PortProtocol: port.Type,
186 NetworkMode: cntr.HostConfig.NetworkMode,
187 NetworkDriver: netDriver,
188 IPAddress: network.IPAddress,
189 }
190 tgt.Address = net.JoinHostPort(tgt.IPAddress, tgt.PrivatePort)
191
192 hash, err := model.CalcHash(tgt)
193 if err != nil {
194 continue
195 }
196
197 tgt.hash = hash
198
199 tgg.targets = append(tgg.targets, tgt)
200 }
201 }
202
203 return tgg
204 }
205
206 func (d *Discoverer) cleanup() {
207 if d.dockerClient != nil {
208 _ = d.dockerClient.Close()
209 }
210 }
211
212 func cntrSource(cntr typesContainer.Summary) string {
213 name := strings.TrimPrefix(cntr.Names[0], "/")
214 return fmt.Sprintf("discoverer=docker,container=%s,image=%s", name, cntr.Image)
215 }