telegraf/plugins/inputs/mongodb/mongodb.go

200 lines
4.8 KiB
Go
Raw Normal View History

2015-07-07 01:20:11 +00:00
package mongodb
import (
"crypto/tls"
"crypto/x509"
2015-07-07 01:20:11 +00:00
"fmt"
"net"
2015-07-07 01:20:11 +00:00
"net/url"
"strings"
2015-07-07 01:20:11 +00:00
"sync"
"time"
"github.com/influxdata/telegraf"
2018-05-04 23:33:23 +00:00
tlsint "github.com/influxdata/telegraf/internal/tls"
2016-01-20 18:57:35 +00:00
"github.com/influxdata/telegraf/plugins/inputs"
2015-07-07 01:20:11 +00:00
"gopkg.in/mgo.v2"
)
type MongoDB struct {
Servers []string
Ssl Ssl
mongos map[string]*Server
GatherClusterStatus bool
GatherPerdbStats bool
GatherColStats bool
ColStatsDbs []string
2018-05-04 23:33:23 +00:00
tlsint.ClientConfig
Log telegraf.Logger
2015-07-07 01:20:11 +00:00
}
type Ssl struct {
Enabled bool
CaCerts []string `toml:"cacerts"`
}
2015-07-07 01:20:11 +00:00
var sampleConfig = `
## An array of URLs of the form:
## "mongodb://" [user ":" pass "@"] host [ ":" port]
## For example:
## mongodb://user:auth_key@10.10.3.30:27017,
## mongodb://10.10.3.33:18832,
servers = ["mongodb://127.0.0.1:27017"]
2018-04-03 23:58:56 +00:00
## When true, collect cluster status
## Note that the query that counts jumbo chunks triggers a COLLSCAN, which
## may have an impact on performance.
# gather_cluster_status = true
2018-04-03 23:58:56 +00:00
## When true, collect per database stats
# gather_perdb_stats = false
## When true, collect per collection stats
# gather_col_stats = false
## List of db where collections stats are collected
## If empty, all db are concerned
# col_stats_dbs = ["local"]
2018-05-04 23:33:23 +00:00
## Optional TLS Config
# tls_ca = "/etc/telegraf/ca.pem"
# tls_cert = "/etc/telegraf/cert.pem"
# tls_key = "/etc/telegraf/key.pem"
## Use TLS but skip chain & host verification
# insecure_skip_verify = false
2015-08-26 15:21:39 +00:00
`
2015-07-07 01:20:11 +00:00
func (m *MongoDB) SampleConfig() string {
return sampleConfig
}
func (*MongoDB) Description() string {
return "Read metrics from one or many MongoDB servers"
}
var localhost = &url.URL{Host: "mongodb://127.0.0.1:27017"}
2015-07-07 01:20:11 +00:00
// Reads stats from all configured servers accumulates stats.
// Returns one of the errors encountered while gather stats (if any).
func (m *MongoDB) Gather(acc telegraf.Accumulator) error {
2015-07-07 01:20:11 +00:00
if len(m.Servers) == 0 {
m.gatherServer(m.getMongoServer(localhost), acc)
return nil
}
var wg sync.WaitGroup
for i, serv := range m.Servers {
if !strings.HasPrefix(serv, "mongodb://") {
// Preserve backwards compatibility for hostnames without a
// scheme, broken in go 1.8. Remove in Telegraf 2.0
serv = "mongodb://" + serv
m.Log.Warnf("Using %q as connection URL; please update your configuration to use an URL", serv)
m.Servers[i] = serv
}
2015-07-07 01:20:11 +00:00
u, err := url.Parse(serv)
if err != nil {
m.Log.Errorf("Unable to parse address %q: %s", serv, err.Error())
2017-04-24 18:13:26 +00:00
continue
2015-07-07 01:20:11 +00:00
}
if u.Host == "" {
m.Log.Errorf("Unable to parse address %q", serv)
continue
}
2015-07-07 01:20:11 +00:00
wg.Add(1)
go func(srv *Server) {
2015-07-07 01:20:11 +00:00
defer wg.Done()
2019-11-12 21:44:57 +00:00
err := m.gatherServer(srv, acc)
if err != nil {
m.Log.Errorf("Error in plugin: %v", err)
}
}(m.getMongoServer(u))
2015-07-07 01:20:11 +00:00
}
wg.Wait()
2017-04-24 18:13:26 +00:00
return nil
2015-07-07 01:20:11 +00:00
}
func (m *MongoDB) getMongoServer(url *url.URL) *Server {
if _, ok := m.mongos[url.Host]; !ok {
m.mongos[url.Host] = &Server{
Log: m.Log,
2015-07-07 01:20:11 +00:00
Url: url,
}
}
return m.mongos[url.Host]
}
func (m *MongoDB) gatherServer(server *Server, acc telegraf.Accumulator) error {
2015-07-07 01:20:11 +00:00
if server.Session == nil {
var dialAddrs []string
if server.Url.User != nil {
dialAddrs = []string{server.Url.String()}
} else {
dialAddrs = []string{server.Url.Host}
}
dialInfo, err := mgo.ParseURL(dialAddrs[0])
if err != nil {
return fmt.Errorf("unable to parse URL %q: %s", dialAddrs[0], err.Error())
2015-07-07 01:20:11 +00:00
}
dialInfo.Direct = true
dialInfo.Timeout = 5 * time.Second
var tlsConfig *tls.Config
if m.Ssl.Enabled {
2018-05-04 23:33:23 +00:00
// Deprecated TLS config
tlsConfig = &tls.Config{}
if len(m.Ssl.CaCerts) > 0 {
roots := x509.NewCertPool()
for _, caCert := range m.Ssl.CaCerts {
ok := roots.AppendCertsFromPEM([]byte(caCert))
if !ok {
return fmt.Errorf("failed to parse root certificate")
}
}
tlsConfig.RootCAs = roots
} else {
tlsConfig.InsecureSkipVerify = true
}
} else {
2018-05-04 23:33:23 +00:00
tlsConfig, err = m.ClientConfig.TLSConfig()
if err != nil {
return err
}
}
// If configured to use TLS, add a dial function
if tlsConfig != nil {
dialInfo.DialServer = func(addr *mgo.ServerAddr) (net.Conn, error) {
conn, err := tls.Dial("tcp", addr.String(), tlsConfig)
if err != nil {
fmt.Printf("error in Dial, %s\n", err.Error())
}
return conn, err
}
}
2015-07-07 01:20:11 +00:00
sess, err := mgo.DialWithInfo(dialInfo)
if err != nil {
return fmt.Errorf("unable to connect to MongoDB: %s", err.Error())
2015-07-07 01:20:11 +00:00
}
server.Session = sess
}
return server.gatherData(acc, m.GatherClusterStatus, m.GatherPerdbStats, m.GatherColStats, m.ColStatsDbs)
2015-07-07 01:20:11 +00:00
}
func init() {
inputs.Add("mongodb", func() telegraf.Input {
2015-07-07 01:20:11 +00:00
return &MongoDB{
mongos: make(map[string]*Server),
GatherClusterStatus: true,
GatherPerdbStats: false,
GatherColStats: false,
ColStatsDbs: []string{"local"},
2015-07-07 01:20:11 +00:00
}
})
}