2021-08-31 10:43:16 +02:00
|
|
|
package main
|
|
|
|
|
|
|
|
import (
|
2021-09-08 09:08:51 +02:00
|
|
|
"context"
|
2022-03-08 09:31:16 +01:00
|
|
|
"errors"
|
2022-01-20 10:14:28 +01:00
|
|
|
"fmt"
|
2021-08-31 10:43:16 +02:00
|
|
|
"log"
|
2022-03-08 09:31:16 +01:00
|
|
|
"net"
|
2021-09-07 09:28:41 +02:00
|
|
|
"sync"
|
2022-01-20 10:14:28 +01:00
|
|
|
"time"
|
2021-08-31 10:43:16 +02:00
|
|
|
|
2021-10-07 14:52:16 +02:00
|
|
|
"github.com/influxdata/line-protocol/v2/lineprotocol"
|
2021-08-31 10:43:16 +02:00
|
|
|
"github.com/nats-io/nats.go"
|
|
|
|
)
|
|
|
|
|
|
|
|
type Metric struct {
|
|
|
|
Name string
|
|
|
|
Value Float
|
2022-03-08 09:31:16 +01:00
|
|
|
|
|
|
|
mc MetricConfig
|
|
|
|
}
|
|
|
|
|
|
|
|
// Currently unused, could be used to send messages via raw TCP.
|
|
|
|
// Each connection is handled in it's own goroutine. This is a blocking function.
|
|
|
|
func ReceiveRaw(ctx context.Context, listener net.Listener, handleLine func(*lineprotocol.Decoder, string) error) error {
|
|
|
|
var wg sync.WaitGroup
|
|
|
|
|
|
|
|
wg.Add(1)
|
|
|
|
go func() {
|
|
|
|
defer wg.Done()
|
|
|
|
<-ctx.Done()
|
|
|
|
if err := listener.Close(); err != nil {
|
|
|
|
log.Printf("listener.Close(): %s", err.Error())
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
for {
|
|
|
|
conn, err := listener.Accept()
|
|
|
|
if err != nil {
|
|
|
|
if errors.Is(err, net.ErrClosed) {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Printf("listener.Accept(): %s", err.Error())
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Add(2)
|
|
|
|
go func() {
|
|
|
|
defer wg.Done()
|
|
|
|
defer conn.Close()
|
|
|
|
|
|
|
|
dec := lineprotocol.NewDecoder(conn)
|
|
|
|
connctx, cancel := context.WithCancel(context.Background())
|
|
|
|
defer cancel()
|
|
|
|
go func() {
|
|
|
|
defer wg.Done()
|
|
|
|
select {
|
|
|
|
case <-connctx.Done():
|
|
|
|
conn.Close()
|
|
|
|
case <-ctx.Done():
|
|
|
|
conn.Close()
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
|
|
|
|
if err := handleLine(dec, "default"); err != nil {
|
|
|
|
if errors.Is(err, net.ErrClosed) {
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
|
|
|
log.Printf("%s: %s", conn.RemoteAddr().String(), err.Error())
|
|
|
|
errmsg := make([]byte, 128)
|
|
|
|
errmsg = append(errmsg, `error: `...)
|
|
|
|
errmsg = append(errmsg, err.Error()...)
|
|
|
|
errmsg = append(errmsg, '\n')
|
|
|
|
conn.Write(errmsg)
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
}
|
|
|
|
|
|
|
|
wg.Wait()
|
|
|
|
return nil
|
2021-08-31 10:43:16 +02:00
|
|
|
}
|
|
|
|
|
|
|
|
// Connect to a nats server and subscribe to "updates". This is a blocking
|
|
|
|
// function. handleLine will be called for each line recieved via nats.
|
|
|
|
// Send `true` through the done channel for gracefull termination.
|
2022-02-22 13:58:10 +01:00
|
|
|
func ReceiveNats(conf *NatsConfig, handleLine func(*lineprotocol.Decoder, string) error, workers int, ctx context.Context) error {
|
2022-02-04 08:52:53 +01:00
|
|
|
var opts []nats.Option
|
|
|
|
if conf.Username != "" && conf.Password != "" {
|
|
|
|
opts = append(opts, nats.UserInfo(conf.Username, conf.Password))
|
|
|
|
}
|
|
|
|
|
|
|
|
nc, err := nats.Connect(conf.Address, opts...)
|
2021-08-31 10:43:16 +02:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
defer nc.Close()
|
|
|
|
|
2021-10-07 14:52:16 +02:00
|
|
|
var wg sync.WaitGroup
|
2022-02-22 14:03:45 +01:00
|
|
|
var subs []*nats.Subscription
|
2021-09-07 09:28:41 +02:00
|
|
|
|
2021-10-07 14:52:16 +02:00
|
|
|
msgs := make(chan *nats.Msg, workers*2)
|
2021-09-08 09:08:51 +02:00
|
|
|
|
2022-02-22 14:03:45 +01:00
|
|
|
for _, sc := range conf.Subscriptions {
|
|
|
|
clusterTag := sc.ClusterTag
|
|
|
|
var sub *nats.Subscription
|
|
|
|
if workers > 1 {
|
|
|
|
wg.Add(workers)
|
|
|
|
|
|
|
|
for i := 0; i < workers; i++ {
|
|
|
|
go func() {
|
|
|
|
for m := range msgs {
|
|
|
|
dec := lineprotocol.NewDecoderWithBytes(m.Data)
|
|
|
|
if err := handleLine(dec, clusterTag); err != nil {
|
|
|
|
log.Printf("error: %s\n", err.Error())
|
|
|
|
}
|
2021-10-11 16:28:05 +02:00
|
|
|
}
|
2021-09-08 09:08:51 +02:00
|
|
|
|
2022-02-22 14:03:45 +01:00
|
|
|
wg.Done()
|
|
|
|
}()
|
2021-10-11 16:28:05 +02:00
|
|
|
}
|
2021-09-08 09:08:51 +02:00
|
|
|
|
2022-02-22 14:03:45 +01:00
|
|
|
sub, err = nc.Subscribe(sc.SubscribeTo, func(m *nats.Msg) {
|
|
|
|
msgs <- m
|
|
|
|
})
|
|
|
|
} else {
|
|
|
|
sub, err = nc.Subscribe(sc.SubscribeTo, func(m *nats.Msg) {
|
|
|
|
dec := lineprotocol.NewDecoderWithBytes(m.Data)
|
|
|
|
if err := handleLine(dec, clusterTag); err != nil {
|
|
|
|
log.Printf("error: %s\n", err.Error())
|
|
|
|
}
|
|
|
|
})
|
|
|
|
}
|
2021-09-07 09:28:41 +02:00
|
|
|
|
2022-02-22 14:03:45 +01:00
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
log.Printf("NATS subscription to '%s' on '%s' established\n", sc.SubscribeTo, conf.Address)
|
|
|
|
subs = append(subs, sub)
|
|
|
|
}
|
2021-10-07 14:52:16 +02:00
|
|
|
|
|
|
|
<-ctx.Done()
|
2022-02-22 14:03:45 +01:00
|
|
|
for _, sub := range subs {
|
|
|
|
err = sub.Unsubscribe()
|
|
|
|
if err != nil {
|
|
|
|
log.Printf("NATS unsubscribe failed: %s", err.Error())
|
|
|
|
}
|
|
|
|
}
|
2021-10-07 14:52:16 +02:00
|
|
|
close(msgs)
|
|
|
|
wg.Wait()
|
|
|
|
|
2021-09-08 09:08:51 +02:00
|
|
|
nc.Close()
|
|
|
|
log.Println("NATS connection closed")
|
|
|
|
return nil
|
2021-08-31 10:43:16 +02:00
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
|
2022-01-24 09:50:12 +01:00
|
|
|
// Place `prefix` in front of `buf` but if possible,
|
|
|
|
// do that inplace in `buf`.
|
|
|
|
func reorder(buf, prefix []byte) []byte {
|
|
|
|
n := len(prefix)
|
|
|
|
m := len(buf)
|
|
|
|
if cap(buf) < m+n {
|
|
|
|
return append(prefix[:n:n], buf...)
|
|
|
|
} else {
|
|
|
|
buf = buf[:n+m]
|
|
|
|
for i := m - 1; i >= 0; i-- {
|
|
|
|
buf[i+n] = buf[i]
|
|
|
|
}
|
|
|
|
for i := 0; i < n; i++ {
|
|
|
|
buf[i] = prefix[i]
|
|
|
|
}
|
|
|
|
return buf
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-02-22 13:58:10 +01:00
|
|
|
// Decode lines using dec and make write calls to the MemoryStore.
|
|
|
|
// If a line is missing its cluster tag, use clusterDefault as default.
|
|
|
|
func decodeLine(dec *lineprotocol.Decoder, clusterDefault string) error {
|
2022-01-21 10:47:40 +01:00
|
|
|
// Reduce allocations in loop:
|
|
|
|
t := time.Now()
|
2022-03-08 09:27:44 +01:00
|
|
|
metric, metricBuf := Metric{}, make([]byte, 0, 16)
|
2022-01-21 10:47:40 +01:00
|
|
|
selector := make([]string, 0, 4)
|
2022-03-08 09:27:44 +01:00
|
|
|
typeBuf, subTypeBuf := make([]byte, 0, 16), make([]byte, 0)
|
2022-01-21 10:47:40 +01:00
|
|
|
|
|
|
|
// Optimize for the case where all lines in a "batch" are about the same
|
|
|
|
// cluster and host. By using `WriteToLevel` (level = host), we do not need
|
|
|
|
// to take the root- and cluster-level lock as often.
|
2022-01-24 09:50:12 +01:00
|
|
|
var lvl *level = nil
|
2022-01-21 10:47:40 +01:00
|
|
|
var prevCluster, prevHost string = "", ""
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
var ok bool
|
2022-01-20 10:14:28 +01:00
|
|
|
for dec.Next() {
|
|
|
|
rawmeasurement, err := dec.Measurement()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
// Needs to be copied because another call to dec.* would
|
|
|
|
// invalidate the returned slice.
|
|
|
|
metricBuf = append(metricBuf[:0], rawmeasurement...)
|
2022-01-24 09:50:12 +01:00
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
// The go compiler optimizes map[string(byteslice)] lookups:
|
|
|
|
metric.mc, ok = memoryStore.metrics[string(rawmeasurement)]
|
|
|
|
if !ok {
|
|
|
|
continue
|
2022-01-24 09:50:12 +01:00
|
|
|
}
|
|
|
|
|
|
|
|
typeBuf, subTypeBuf := typeBuf[:0], subTypeBuf[:0]
|
2022-02-22 13:58:10 +01:00
|
|
|
cluster, host := clusterDefault, ""
|
2022-01-20 10:14:28 +01:00
|
|
|
for {
|
|
|
|
key, val, err := dec.NextTag()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
|
|
|
}
|
|
|
|
if key == nil {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
2022-01-21 10:47:40 +01:00
|
|
|
// The go compiler optimizes string([]byte{...}) == "...":
|
2022-01-20 10:14:28 +01:00
|
|
|
switch string(key) {
|
|
|
|
case "cluster":
|
2022-01-21 10:47:40 +01:00
|
|
|
if string(val) == prevCluster {
|
|
|
|
cluster = prevCluster
|
|
|
|
} else {
|
|
|
|
cluster = string(val)
|
2022-01-24 09:50:12 +01:00
|
|
|
lvl = nil
|
2022-01-21 10:47:40 +01:00
|
|
|
}
|
2022-01-24 09:50:12 +01:00
|
|
|
case "hostname", "host":
|
2022-01-21 10:47:40 +01:00
|
|
|
if string(val) == prevHost {
|
|
|
|
host = prevHost
|
|
|
|
} else {
|
|
|
|
host = string(val)
|
2022-01-24 09:50:12 +01:00
|
|
|
lvl = nil
|
2022-01-21 10:47:40 +01:00
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
case "type":
|
2022-01-24 09:50:12 +01:00
|
|
|
if string(val) == "node" {
|
|
|
|
break
|
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
// We cannot be sure that the "type" tag comes before the "type-id" tag:
|
2022-01-24 09:50:12 +01:00
|
|
|
if len(typeBuf) == 0 {
|
|
|
|
typeBuf = append(typeBuf, val...)
|
|
|
|
} else {
|
|
|
|
typeBuf = reorder(typeBuf, val)
|
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
case "type-id":
|
2022-01-24 09:50:12 +01:00
|
|
|
typeBuf = append(typeBuf, val...)
|
2022-01-20 10:14:28 +01:00
|
|
|
case "subtype":
|
2022-03-08 09:27:44 +01:00
|
|
|
// We cannot be sure that the "subtype" tag comes before the "stype-id" tag:
|
2022-01-24 09:50:12 +01:00
|
|
|
if len(subTypeBuf) == 0 {
|
|
|
|
subTypeBuf = append(subTypeBuf, val...)
|
|
|
|
} else {
|
|
|
|
subTypeBuf = reorder(typeBuf, val)
|
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
case "stype-id":
|
2022-01-24 09:50:12 +01:00
|
|
|
subTypeBuf = append(subTypeBuf, val...)
|
2022-01-20 10:14:28 +01:00
|
|
|
default:
|
|
|
|
// Ignore unkown tags (cc-metric-collector might send us a unit for example that we do not need)
|
|
|
|
// return fmt.Errorf("unkown tag: '%s' (value: '%s')", string(key), string(val))
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
// If the cluster or host changed, the lvl was set to nil
|
2022-01-24 09:50:12 +01:00
|
|
|
if lvl == nil {
|
2022-01-21 10:47:40 +01:00
|
|
|
selector = selector[:2]
|
2022-01-24 09:50:12 +01:00
|
|
|
selector[0], selector[1] = cluster, host
|
|
|
|
lvl = memoryStore.GetLevel(selector)
|
|
|
|
prevCluster, prevHost = cluster, host
|
2022-01-21 10:47:40 +01:00
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
// subtypes:
|
2022-01-21 10:47:40 +01:00
|
|
|
selector = selector[:0]
|
2022-01-24 09:50:12 +01:00
|
|
|
if len(typeBuf) > 0 {
|
2022-01-21 10:47:40 +01:00
|
|
|
selector = append(selector, string(typeBuf)) // <- Allocation :(
|
2022-01-24 09:50:12 +01:00
|
|
|
if len(subTypeBuf) > 0 {
|
2022-01-21 10:47:40 +01:00
|
|
|
selector = append(selector, string(subTypeBuf))
|
2022-01-20 10:14:28 +01:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
for {
|
|
|
|
key, val, err := dec.NextField()
|
|
|
|
if err != nil {
|
|
|
|
return err
|
2022-01-20 10:14:28 +01:00
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
if key == nil {
|
|
|
|
break
|
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
if string(key) != "value" {
|
|
|
|
return fmt.Errorf("unkown field: '%s' (value: %#v)", string(key), val)
|
2022-01-20 10:14:28 +01:00
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
if val.Kind() == lineprotocol.Float {
|
|
|
|
metric.Value = Float(val.FloatV())
|
|
|
|
} else if val.Kind() == lineprotocol.Int {
|
|
|
|
metric.Value = Float(val.IntV())
|
|
|
|
} else {
|
|
|
|
return fmt.Errorf("unsupported value type in message: %s", val.Kind().String())
|
|
|
|
}
|
2022-01-20 10:14:28 +01:00
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
if t, err = dec.Time(lineprotocol.Second, t); err != nil {
|
2022-01-20 10:14:28 +01:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
|
2022-03-08 09:27:44 +01:00
|
|
|
if err := memoryStore.WriteToLevel(lvl, selector, t.Unix(), []Metric{metric}); err != nil {
|
2022-01-20 10:14:28 +01:00
|
|
|
return err
|
|
|
|
}
|
|
|
|
}
|
|
|
|
return nil
|
|
|
|
}
|