/
mqtt.go
75 lines (64 loc) · 2.43 KB
/
mqtt.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
// SPDX-License-Identifier: MIT OR Apache-2.0
package outputs
import (
"crypto/tls"
"log"
"github.com/DataDog/datadog-go/statsd"
mqtt "github.com/eclipse/paho.mqtt.golang"
"github.com/google/uuid"
"github.com/falcosecurity/falcosidekick/types"
)
// NewMQTTClient returns a new output.Client for accessing Kubernetes.
func NewMQTTClient(config *types.Configuration, stats *types.Statistics, promStats *types.PromStatistics,
statsdClient, dogstatsdClient *statsd.Client) (*Client, error) {
options := mqtt.NewClientOptions()
options.AddBroker(config.MQTT.Broker)
options.SetClientID("falcosidekick-" + uuid.NewString()[:6])
if config.MQTT.User != "" && config.MQTT.Password != "" {
options.Username = config.MQTT.User
options.Password = config.MQTT.Password
}
if !config.MQTT.CheckCert {
options.TLSConfig = &tls.Config{
InsecureSkipVerify: true, // #nosec G402 This is only set as a result of explicit configuration
}
}
options.OnConnectionLost = func(client mqtt.Client, err error) {
log.Printf("[ERROR] : MQTT - Connection lost: %v\n", err.Error())
}
client := mqtt.NewClient(options)
return &Client{
OutputType: MQTT,
Config: config,
MQTTClient: client,
Stats: stats,
PromStats: promStats,
StatsdClient: statsdClient,
DogstatsdClient: dogstatsdClient,
}, nil
}
// MQTTPublish .
func (c *Client) MQTTPublish(falcopayload types.FalcoPayload) {
c.Stats.MQTT.Add(Total, 1)
t := c.MQTTClient.Connect()
t.Wait()
if err := t.Error(); err != nil {
go c.CountMetric(Outputs, 1, []string{"output:mqtt", "status:error"})
c.Stats.MQTT.Add(Error, 1)
c.PromStats.Outputs.With(map[string]string{"destination": "mqtt", "status": err.Error()}).Inc()
log.Printf("[ERROR] : %s - %v\n", MQTT, err.Error())
return
}
defer c.MQTTClient.Disconnect(100)
if err := c.MQTTClient.Publish(c.Config.MQTT.Topic, byte(c.Config.MQTT.QOS), c.Config.MQTT.Retained, falcopayload.String()).Error(); err != nil {
go c.CountMetric(Outputs, 1, []string{"output:mqtt", "status:error"})
c.Stats.MQTT.Add(Error, 1)
c.PromStats.Outputs.With(map[string]string{"destination": "mqtt", "status": Error}).Inc()
log.Printf("[ERROR] : %s - %v\n", MQTT, err.Error())
return
}
log.Printf("[INFO] : %s - Message published\n", MQTT)
go c.CountMetric(Outputs, 1, []string{"output:mqtt", "status:ok"})
c.Stats.MQTT.Add(OK, 1)
c.PromStats.Outputs.With(map[string]string{"destination": "mqtt", "status": OK}).Inc()
}