mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-01 15:07:17 +08:00
lbagent: add telegraf support
This commit is contained in:
@@ -5,4 +5,5 @@ REQUIRES=(
|
||||
"keepalived >= 2.0.0"
|
||||
"haproxy >= 1.8.0"
|
||||
"gobetween >= 0.7.0"
|
||||
"telegraf >= 1.5"
|
||||
)
|
||||
|
||||
@@ -18,6 +18,7 @@ func init() {
|
||||
keys := []string{
|
||||
"params.keepalived_conf_tmpl",
|
||||
"params.haproxy_conf_tmpl",
|
||||
"params.telegraf_conf_tmpl",
|
||||
}
|
||||
d := data.(*jsonutils.JSONDict)
|
||||
for _, key := range keys {
|
||||
|
||||
@@ -3,6 +3,7 @@ package models
|
||||
import (
|
||||
"context"
|
||||
"encoding/base64"
|
||||
"net/url"
|
||||
"reflect"
|
||||
"text/template"
|
||||
"time"
|
||||
@@ -78,11 +79,19 @@ type SLoadbalancerAgentParamsHaproxy struct {
|
||||
LogNormal bool
|
||||
}
|
||||
|
||||
type SLoadbalancerAgentParamsTelegraf struct {
|
||||
InfluxDbOutputUrl string
|
||||
InfluxDbOutputName string
|
||||
HaproxyInputInterval int
|
||||
}
|
||||
|
||||
type SLoadbalancerAgentParams struct {
|
||||
KeepalivedConfTmpl string
|
||||
HaproxyConfTmpl string
|
||||
TelegrafConfTmpl string
|
||||
Vrrp SLoadbalancerAgentParamsVrrp
|
||||
Haproxy SLoadbalancerAgentParamsHaproxy
|
||||
Telegraf SLoadbalancerAgentParamsTelegraf
|
||||
}
|
||||
|
||||
func (p *SLoadbalancerAgentParamsVrrp) Validate(data *jsonutils.JSONDict) error {
|
||||
@@ -141,6 +150,31 @@ func (p *SLoadbalancerAgentParamsHaproxy) initDefault(data *jsonutils.JSONDict)
|
||||
}
|
||||
}
|
||||
|
||||
func (p *SLoadbalancerAgentParamsTelegraf) Validate(data *jsonutils.JSONDict) error {
|
||||
if p.InfluxDbOutputUrl != "" {
|
||||
_, err := url.Parse(p.InfluxDbOutputUrl)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
if p.HaproxyInputInterval <= 0 {
|
||||
p.HaproxyInputInterval = 5
|
||||
}
|
||||
if p.InfluxDbOutputName == "" {
|
||||
p.InfluxDbOutputName = "telegraf"
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p *SLoadbalancerAgentParamsTelegraf) initDefault(data *jsonutils.JSONDict) {
|
||||
if p.HaproxyInputInterval == 0 {
|
||||
p.HaproxyInputInterval = 5
|
||||
}
|
||||
if p.InfluxDbOutputName == "" {
|
||||
p.InfluxDbOutputName = "telegraf"
|
||||
}
|
||||
}
|
||||
|
||||
func (p *SLoadbalancerAgentParams) validateTmpl(k, s string) error {
|
||||
d, err := base64.StdEncoding.DecodeString(s)
|
||||
if err != nil {
|
||||
@@ -161,8 +195,12 @@ func (p *SLoadbalancerAgentParams) initDefault(data *jsonutils.JSONDict) {
|
||||
if p.HaproxyConfTmpl == "" {
|
||||
p.HaproxyConfTmpl = loadbalancerHaproxyConfTmplDefaultEncoded
|
||||
}
|
||||
if p.TelegrafConfTmpl == "" {
|
||||
p.TelegrafConfTmpl = loadbalancerTelegrafConfTmplDefaultEncoded
|
||||
}
|
||||
p.Vrrp.initDefault(data)
|
||||
p.Haproxy.initDefault(data)
|
||||
p.Telegraf.initDefault(data)
|
||||
}
|
||||
|
||||
func (p *SLoadbalancerAgentParams) Validate(data *jsonutils.JSONDict) error {
|
||||
@@ -173,12 +211,18 @@ func (p *SLoadbalancerAgentParams) Validate(data *jsonutils.JSONDict) error {
|
||||
if err := p.validateTmpl("haproxy_conf_tmpl", p.HaproxyConfTmpl); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := p.validateTmpl("telegraf_conf_tmpl", p.TelegrafConfTmpl); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := p.Vrrp.Validate(data); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := p.Haproxy.Validate(data); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := p.Telegraf.Validate(data); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -459,9 +503,21 @@ listen stats
|
||||
stats auth Yunion:LBStats
|
||||
stats uri /
|
||||
`
|
||||
|
||||
loadbalancerTelegrafConfTmplDefault = `
|
||||
[[outputs.influxdb]]
|
||||
urls = ["{{ .telegraf.influx_db_output_url }}"]
|
||||
database = "{{ .telegraf.influx_db_output_name }}"
|
||||
|
||||
[[inputs.haproxy]]
|
||||
interval = "{{ .telegraf.haproxy_input_interval }}s"
|
||||
servers = ["{{ .telegraf.haproxy_input_stats_socket }}"]
|
||||
keep_field_names = true
|
||||
`
|
||||
)
|
||||
|
||||
var (
|
||||
loadbalancerKeepalivedConfTmplDefaultEncoded = base64.StdEncoding.EncodeToString([]byte(loadbalancerKeepalivedConfTmplDefault))
|
||||
loadbalancerHaproxyConfTmplDefaultEncoded = base64.StdEncoding.EncodeToString([]byte(loadbalancerHaproxyConfTmplDefault))
|
||||
loadbalancerTelegrafConfTmplDefaultEncoded = base64.StdEncoding.EncodeToString([]byte(loadbalancerTelegrafConfTmplDefault))
|
||||
)
|
||||
|
||||
@@ -1,8 +1,10 @@
|
||||
package lbagent
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"os/exec"
|
||||
"path/filepath"
|
||||
@@ -142,6 +144,28 @@ func (h *HaproxyHelper) handleUseCorpusCmd(ctx context.Context, cmd *LbagentCmd)
|
||||
return err
|
||||
}
|
||||
}
|
||||
if agentParams.AgentModel.Params.Telegraf.InfluxDbOutputUrl != "" {
|
||||
agentParams.SetTelegrafParams("haproxy_input_stats_socket", h.haproxyStatsSocketFile())
|
||||
// telegraf config
|
||||
buf := bytes.NewBufferString("# yunion lb auto-generated telegraf.conf\n")
|
||||
tmpl := agentParams.TelegrafConfigTmpl
|
||||
err := tmpl.Execute(buf, agentParams.Data)
|
||||
if err == nil {
|
||||
d := buf.Bytes()
|
||||
p := filepath.Join(dir, "telegraf.conf")
|
||||
err := ioutil.WriteFile(p, d, agentutils.FileModeFile)
|
||||
if err == nil {
|
||||
err := h.reloadTelegraf(ctx)
|
||||
if err != nil {
|
||||
log.Errorf("reloading telegraf.conf failed: %s", err)
|
||||
}
|
||||
} else {
|
||||
log.Errorf("writing %s failed: %s", p, err)
|
||||
}
|
||||
} else {
|
||||
log.Errorf("making telegraf.conf failed: %s, tmpl:\n%s", err, tmpl)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
@@ -284,6 +308,38 @@ func (h *HaproxyHelper) reloadGobetween(ctx context.Context) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *HaproxyHelper) telegrafConf() string {
|
||||
return filepath.Join(h.haproxyConfD(), "telegraf.conf")
|
||||
}
|
||||
|
||||
func (h *HaproxyHelper) telegrafPidFile() string {
|
||||
return filepath.Join(h.opts.haproxyRunDir, "telegraf.pid")
|
||||
}
|
||||
|
||||
func (h *HaproxyHelper) reloadTelegraf(ctx context.Context) error {
|
||||
pidFile := h.telegrafPidFile()
|
||||
args := []string{
|
||||
h.opts.TelegrafBin,
|
||||
"--config", h.telegrafConf(),
|
||||
}
|
||||
proc := agentutils.ReadPidFile(pidFile)
|
||||
if proc != nil {
|
||||
log.Infof("stopping telegraf(%d)", proc.Pid)
|
||||
proc.Kill()
|
||||
proc.Wait()
|
||||
}
|
||||
log.Infof("starting telegraf")
|
||||
cmd, err := h.startCmd(args)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
err = agentutils.WritePidFile(cmd.Process.Pid, h.telegrafPidFile())
|
||||
if err != nil {
|
||||
return fmt.Errorf("writing telegraf pid file: %s", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (h *HaproxyHelper) keepalivedConf() string {
|
||||
return filepath.Join(h.opts.haproxyConfigDir, "keepalived.conf")
|
||||
}
|
||||
|
||||
@@ -12,13 +12,15 @@ type AgentParams struct {
|
||||
AgentModel *models.LoadbalancerAgent
|
||||
KeepalivedConfigTmpl *template.Template
|
||||
HaproxyConfigTmpl *template.Template
|
||||
Data map[string]interface{}
|
||||
TelegrafConfigTmpl *template.Template
|
||||
Data map[string]map[string]interface{}
|
||||
}
|
||||
|
||||
func NewAgentParams(agent *models.LoadbalancerAgent) (*AgentParams, error) {
|
||||
b64s := map[string]string{
|
||||
"keepalived_conf_tmpl": agent.Params.KeepalivedConfTmpl,
|
||||
"haproxy_conf_tmpl": agent.Params.HaproxyConfTmpl,
|
||||
"telegraf_conf_tmpl": agent.Params.TelegrafConfTmpl,
|
||||
}
|
||||
tmpls := map[string]*template.Template{}
|
||||
for name, b64 := range b64s {
|
||||
@@ -52,15 +54,22 @@ func NewAgentParams(agent *models.LoadbalancerAgent) (*AgentParams, error) {
|
||||
"log_tcp": agent.Params.Haproxy.LogTcp,
|
||||
"log_normal": agent.Params.Haproxy.LogNormal,
|
||||
}
|
||||
data := map[string]interface{}{
|
||||
"agent": dataAgent,
|
||||
"vrrp": dataVrrp,
|
||||
"haproxy": dataHaproxy,
|
||||
dataTelegraf := map[string]interface{}{
|
||||
"influx_db_output_url": agent.Params.Telegraf.InfluxDbOutputUrl,
|
||||
"influx_db_output_name": agent.Params.Telegraf.InfluxDbOutputName,
|
||||
"haproxy_input_interval": agent.Params.Telegraf.HaproxyInputInterval,
|
||||
}
|
||||
data := map[string]map[string]interface{}{
|
||||
"agent": dataAgent,
|
||||
"vrrp": dataVrrp,
|
||||
"haproxy": dataHaproxy,
|
||||
"telegraf": dataTelegraf,
|
||||
}
|
||||
agentParams := &AgentParams{
|
||||
AgentModel: agent,
|
||||
KeepalivedConfigTmpl: tmpls["keepalived_conf_tmpl"],
|
||||
HaproxyConfigTmpl: tmpls["haproxy_conf_tmpl"],
|
||||
TelegrafConfigTmpl: tmpls["telegraf_conf_tmpl"],
|
||||
Data: data,
|
||||
}
|
||||
return agentParams, nil
|
||||
@@ -88,7 +97,7 @@ func (p *AgentParams) setXxParams(xx, k string, v interface{}) map[string]interf
|
||||
dt = map[string]interface{}{}
|
||||
p.Data[xx] = dt
|
||||
} else {
|
||||
dt = d.(map[string]interface{})
|
||||
dt = d
|
||||
}
|
||||
dt[k] = v
|
||||
return dt
|
||||
@@ -102,5 +111,9 @@ func (p *AgentParams) SetHaproxyParams(k string, v interface{}) map[string]inter
|
||||
return p.setXxParams("haproxy", k, v)
|
||||
}
|
||||
|
||||
func (p *AgentParams) SetTelegrafParams(k string, v interface{}) map[string]interface{} {
|
||||
return p.setXxParams("telegraf", k, v)
|
||||
}
|
||||
|
||||
func (p *AgentParams) KeepalivedConfig() {
|
||||
}
|
||||
|
||||
@@ -27,6 +27,7 @@ type LbagentOptions struct {
|
||||
KeepalivedBin string `default:"keepalived"`
|
||||
HaproxyBin string `default:"haproxy"`
|
||||
GobetweenBin string `default:"gobetween"`
|
||||
TelegrafBin string `default:"telegraf"`
|
||||
}
|
||||
|
||||
type Options struct {
|
||||
|
||||
@@ -171,9 +171,17 @@ type LoadbalancerAgentParamsHaproxy struct {
|
||||
LogNormal bool
|
||||
}
|
||||
|
||||
type LoadbalancerAgentParamsTelegraf struct {
|
||||
InfluxDbOutputUrl string
|
||||
InfluxDbOutputName string
|
||||
HaproxyInputInterval int
|
||||
}
|
||||
|
||||
type LoadbalancerAgentParams struct {
|
||||
KeepalivedConfTmpl string
|
||||
HaproxyConfTmpl string
|
||||
TelegrafConfTmpl string
|
||||
Vrrp LoadbalancerAgentParamsVrrp
|
||||
Haproxy LoadbalancerAgentParamsHaproxy
|
||||
Telegraf LoadbalancerAgentParamsTelegraf
|
||||
}
|
||||
|
||||
@@ -20,10 +20,14 @@ type LoadbalancerAgentParamsOptions struct {
|
||||
VrrpPass string
|
||||
|
||||
HaproxyGlobalLog string
|
||||
HaproxyGlobalNbthread *int `default:"1" help:"enable experimental multi-threading support available since haproxy 1.8"`
|
||||
HaproxyGlobalNbthread *int `help:"enable experimental multi-threading support available since haproxy 1.8"`
|
||||
HaproxyLogHttp string `choices:"true|false"`
|
||||
HaproxyLogTcp string `choices:"true|false"`
|
||||
HaproxyLogNormal string `choices:"true|false"`
|
||||
|
||||
TelegrafInfluxDbOutputUrl string
|
||||
TelegrafInfluxDbOutputName string
|
||||
TelegrafHaproxyInputInterval int
|
||||
}
|
||||
|
||||
func (opts *LoadbalancerAgentParamsOptions) setPrefixedParams(params *jsonutils.JSONDict, pref string) {
|
||||
@@ -49,6 +53,7 @@ func (opts *LoadbalancerAgentParamsOptions) Params() (*jsonutils.JSONDict, error
|
||||
}
|
||||
opts.setPrefixedParams(params, "vrrp")
|
||||
opts.setPrefixedParams(params, "haproxy")
|
||||
opts.setPrefixedParams(params, "telegraf")
|
||||
return params, nil
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user