398 строки
11 KiB
Go
398 строки
11 KiB
Go
// Package influx provides functionality.
|
|
package influx
|
|
|
|
import (
|
|
"bytes"
|
|
"encoding/json"
|
|
"fmt"
|
|
"io"
|
|
"log"
|
|
"net/http"
|
|
"net/url"
|
|
"os"
|
|
"strconv"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/davecgh/go-spew/spew"
|
|
)
|
|
|
|
const (
|
|
defaultAddr = "http://localhost:8428"
|
|
// VictoriaMetrics uses MetricsQL, not Flux
|
|
// Data model: metric names are {measurement}_{field}
|
|
)
|
|
|
|
var (
|
|
initOnce sync.Once
|
|
client *http.Client
|
|
// addr defaults to defaultAddr at declaration so callers (and
|
|
// tests) can override it before the first lazy init runs.
|
|
addr = defaultAddr
|
|
)
|
|
|
|
// ensureInit performs one-time setup of the HTTP client and TSDB
|
|
// address. It is called from every public function so that binaries
|
|
// which only import this package transitively (e.g. the distributed
|
|
// worker, which never reads or writes TSDB points) do not pay the
|
|
// init cost or emit a misleading "TSDB client initialized" log line.
|
|
// INFLUX_URL is read here so the address tracks the env var across
|
|
// process restarts without requiring explicit init from the caller.
|
|
//
|
|
// The env var is honored only when `addr` still equals the default.
|
|
// This lets tests (and any explicit caller) pre-set `addr` to a
|
|
// mock URL before the first public call fires; without this guard,
|
|
// CI runs where INFLUX_URL is exported in .env.ci.example would
|
|
// overwrite a test's mock server URL the moment initOnce fires,
|
|
// causing every QueryVM/* test to silently target a real TSDB.
|
|
func ensureInit() {
|
|
initOnce.Do(func() {
|
|
client = &http.Client{Timeout: 30 * time.Second}
|
|
if addr == defaultAddr {
|
|
if addrEnv := os.Getenv("INFLUX_URL"); addrEnv != "" {
|
|
addr = addrEnv
|
|
}
|
|
}
|
|
log.Println("TSDB client initialized for:", addr)
|
|
})
|
|
}
|
|
|
|
// VMExportResponse represents VictoriaMetrics export response
|
|
type VMExportResponse struct {
|
|
Metric map[string]string `json:"metric"`
|
|
Values []json.Number `json:"values"`
|
|
Timestamps []int64 `json:"timestamps"`
|
|
}
|
|
|
|
// QueryVM performs a MetricsQL query against VictoriaMetrics
|
|
// QueryVM performs an export query against VictoriaMetrics.
|
|
// selector is a time series selector like `chttp_took{check="1146"}`.
|
|
// start is the RFC3339 or Unix timestamp for the beginning of the time range (can be empty).
|
|
func QueryVM(selector, start string) ([]VMExportResponse, error) {
|
|
return QueryVMMany([]string{selector}, start)
|
|
}
|
|
|
|
// QueryVMMany exports several selectors in one request.
|
|
func QueryVMMany(selectors []string, start string) ([]VMExportResponse, error) {
|
|
ensureInit()
|
|
u, err := url.Parse(addr + "/api/v1/export")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
req, err := http.NewRequest("GET", u.String(), http.NoBody)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
q := req.URL.Query()
|
|
for _, selector := range selectors {
|
|
q.Add("match[]", selector)
|
|
}
|
|
if start != "" {
|
|
q.Add("start", start)
|
|
}
|
|
req.URL.RawQuery = q.Encode()
|
|
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
defer resp.Body.Close() //nolint:errcheck
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
body, _ := io.ReadAll(resp.Body)
|
|
return nil, fmt.Errorf("query failed with status %d: %s", resp.StatusCode, string(body))
|
|
}
|
|
|
|
var results []VMExportResponse
|
|
decoder := json.NewDecoder(resp.Body)
|
|
for {
|
|
var result VMExportResponse
|
|
if err := decoder.Decode(&result); err != nil {
|
|
if err == io.EOF {
|
|
break
|
|
}
|
|
return nil, err
|
|
}
|
|
results = append(results, result)
|
|
}
|
|
|
|
return results, nil
|
|
}
|
|
|
|
// MetricCheck identifies the duration series for a check.
|
|
type MetricCheck struct {
|
|
Metric string
|
|
CheckID int64
|
|
}
|
|
|
|
// GetLastMany fetches a page's checks in one bounded VictoriaMetrics export.
|
|
func GetLastMany(checks []MetricCheck, hours int) (map[int64][]InfluxData, error) {
|
|
out := make(map[int64][]InfluxData, len(checks))
|
|
if len(checks) == 0 {
|
|
return out, nil
|
|
}
|
|
if len(checks) > 500 {
|
|
checks = checks[:500]
|
|
}
|
|
selectors := make([]string, 0, len(checks))
|
|
for _, check := range checks {
|
|
selectors = append(selectors, fmt.Sprintf(`%s_took{check="%d"}`, check.Metric, check.CheckID))
|
|
}
|
|
results, err := QueryVMMany(selectors, fmt.Sprintf("%d", time.Now().Add(-time.Duration(hours)*time.Hour).Unix()))
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
for _, result := range results {
|
|
checkID, err := strconv.ParseInt(result.Metric["check"], 10, 64)
|
|
if err != nil {
|
|
continue
|
|
}
|
|
state := result.Metric["state"]
|
|
if state == "" {
|
|
state = "UNK"
|
|
}
|
|
for i, ts := range result.Timestamps {
|
|
var duration int64
|
|
if i < len(result.Values) {
|
|
duration, _ = result.Values[i].Int64()
|
|
}
|
|
out[checkID] = append(out[checkID], InfluxData{Time: time.Unix(ts/1000, (ts%1000)*1e6), Duration: duration, State: state})
|
|
}
|
|
}
|
|
return out, nil
|
|
}
|
|
|
|
// InfluxData provides functionality. //nolint:revive // stutter intentional for clarity
|
|
type InfluxData struct {
|
|
Time time.Time `json:"time"`
|
|
Duration int64 `json:"duration"`
|
|
State string `json:"state"`
|
|
Error string `json:"error,omitempty"`
|
|
Warnings string `json:"warnings,omitempty"`
|
|
}
|
|
|
|
// GetLast provides functionality.
|
|
func GetLast(metric string, check int64, hours int) ([]InfluxData, error) {
|
|
ensureInit()
|
|
// VictoriaMetrics export API: match[] selector + start time
|
|
selector := fmt.Sprintf(`%s_took{check="%d"}`, metric, check)
|
|
start := fmt.Sprintf("%d", time.Now().Add(-time.Duration(hours)*time.Hour).Unix())
|
|
|
|
log.Println("TSDB query:", selector, "start:", start)
|
|
results, err := QueryVM(selector, start)
|
|
if err != nil {
|
|
spew.Dump(err)
|
|
log.Println("TSDB query error", err)
|
|
return nil, err
|
|
}
|
|
|
|
log.Printf("Got %d result series from TSDB", len(results))
|
|
|
|
influxData := make([]InfluxData, 0)
|
|
|
|
// Process each time series (VictoriaMetrics returns one series per unique tag combination)
|
|
for _, result := range results {
|
|
state := result.Metric["state"]
|
|
errorMsg := result.Metric["error"]
|
|
warnings := result.Metric["warnings"]
|
|
|
|
for i, ts := range result.Timestamps {
|
|
// Convert milliseconds to time.Time
|
|
t := time.Unix(ts/1000, (ts%1000)*1e6)
|
|
|
|
// Parse the value
|
|
var duration int64
|
|
if i < len(result.Values) {
|
|
if f, err := result.Values[i].Int64(); err == nil {
|
|
duration = f
|
|
}
|
|
}
|
|
|
|
data := InfluxData{
|
|
Time: t,
|
|
Duration: duration,
|
|
State: state,
|
|
Error: errorMsg,
|
|
Warnings: warnings,
|
|
}
|
|
if data.State == "" {
|
|
data.State = "UNK"
|
|
}
|
|
influxData = append(influxData, data)
|
|
}
|
|
}
|
|
|
|
// Sort by time descending (newest first)
|
|
for i := 0; i < len(influxData); i++ {
|
|
for j := i + 1; j < len(influxData); j++ {
|
|
if influxData[i].Time.Before(influxData[j].Time) {
|
|
influxData[i], influxData[j] = influxData[j], influxData[i]
|
|
}
|
|
}
|
|
}
|
|
|
|
return influxData, nil
|
|
}
|
|
|
|
// escapeTagValue escapes special characters in influx line protocol tag keys/values.
|
|
// Characters that must be escaped: comma, equals, space.
|
|
func escapeTagValue(s string) string {
|
|
s = strings.ReplaceAll(s, `\`, `\\`)
|
|
s = strings.ReplaceAll(s, " ", `\ `)
|
|
s = strings.ReplaceAll(s, ",", `\,`)
|
|
s = strings.ReplaceAll(s, "=", `\=`)
|
|
return s
|
|
}
|
|
|
|
// formatInfluxLine formats data as InfluxDB line protocol
|
|
func formatInfluxLine(measurement string, tags map[string]string, fields map[string]interface{}, ts time.Time) string {
|
|
var buf bytes.Buffer
|
|
|
|
// Write measurement
|
|
buf.WriteString(measurement)
|
|
|
|
// Write tags
|
|
tagKeys := make([]string, 0, len(tags))
|
|
for k := range tags {
|
|
tagKeys = append(tagKeys, k)
|
|
}
|
|
// Sort tags for consistency
|
|
for i := 0; i < len(tagKeys); i++ {
|
|
for j := i + 1; j < len(tagKeys); j++ {
|
|
if tagKeys[i] > tagKeys[j] {
|
|
tagKeys[i], tagKeys[j] = tagKeys[j], tagKeys[i]
|
|
}
|
|
}
|
|
}
|
|
|
|
for _, k := range tagKeys {
|
|
buf.WriteByte(',')
|
|
buf.WriteString(escapeTagValue(k))
|
|
buf.WriteByte('=')
|
|
buf.WriteString(escapeTagValue(tags[k]))
|
|
}
|
|
|
|
buf.WriteByte(' ')
|
|
|
|
// Write fields
|
|
fieldKeys := make([]string, 0, len(fields))
|
|
for k := range fields {
|
|
fieldKeys = append(fieldKeys, k)
|
|
}
|
|
firstField := true
|
|
for _, k := range fieldKeys {
|
|
if !firstField {
|
|
buf.WriteByte(',')
|
|
}
|
|
firstField = false
|
|
buf.WriteString(k)
|
|
buf.WriteByte('=')
|
|
|
|
switch v := fields[k].(type) {
|
|
case int64:
|
|
buf.WriteString(strconv.FormatInt(v, 10) + "i")
|
|
case int:
|
|
buf.WriteString(strconv.FormatInt(int64(v), 10) + "i")
|
|
case float64:
|
|
buf.WriteString(strconv.FormatFloat(v, 'f', -1, 64))
|
|
case bool:
|
|
if v {
|
|
buf.WriteString("true")
|
|
} else {
|
|
buf.WriteString("false")
|
|
}
|
|
case string:
|
|
buf.WriteByte('"')
|
|
buf.WriteString(strings.ReplaceAll(v, "\"", "\\\""))
|
|
buf.WriteByte('"')
|
|
default:
|
|
buf.WriteString(strconv.FormatFloat(0, 'f', -1, 64))
|
|
}
|
|
}
|
|
|
|
// Write timestamp (nanoseconds)
|
|
buf.WriteByte(' ')
|
|
buf.WriteString(strconv.FormatInt(ts.UnixNano(), 10))
|
|
|
|
return buf.String()
|
|
}
|
|
|
|
// WriteOne provides functionality.
|
|
func WriteOne(metric string, tags map[string]string, fields map[string]interface{}) error {
|
|
ensureInit()
|
|
// Format as InfluxDB line protocol
|
|
line := formatInfluxLine(metric, tags, fields, time.Now())
|
|
|
|
// Write to VictoriaMetrics /api/v2/write endpoint
|
|
u, err := url.Parse(addr + "/api/v2/write")
|
|
if err != nil {
|
|
log.Println("TSDB write error (URL parse):", err)
|
|
return err
|
|
}
|
|
|
|
req, err := http.NewRequest("POST", u.String(), bytes.NewBufferString(line))
|
|
if err != nil {
|
|
log.Println("TSDB write error (request):", err)
|
|
return err
|
|
}
|
|
|
|
req.Header.Set("Content-Type", "text/plain")
|
|
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
log.Println("TSDB write error:", err)
|
|
return err
|
|
}
|
|
defer resp.Body.Close() //nolint:errcheck
|
|
|
|
if resp.StatusCode != http.StatusNoContent && resp.StatusCode != http.StatusOK {
|
|
body, _ := io.ReadAll(resp.Body)
|
|
log.Printf("TSDB write failed with status %d: %s\n", resp.StatusCode, string(body))
|
|
return fmt.Errorf("write failed with status %d", resp.StatusCode)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// HealthCheck performs a simple health check against VictoriaMetrics
|
|
func HealthCheck() error {
|
|
ensureInit()
|
|
u, err := url.Parse(addr + "/health")
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
req, err := http.NewRequest("GET", u.String(), http.NoBody)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
resp, err := client.Do(req)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
defer resp.Body.Close() //nolint:errcheck
|
|
|
|
if resp.StatusCode != http.StatusOK {
|
|
return fmt.Errorf("health check failed with status %d", resp.StatusCode)
|
|
}
|
|
|
|
return nil
|
|
}
|
|
|
|
// QueryDB is deprecated and kept for compatibility
|
|
// Use QueryVM for MetricsQL queries instead
|
|
func QueryDB(query string) ([]VMExportResponse, error) {
|
|
// This is a compatibility function for old code
|
|
// Note: Flux queries are NOT supported by VictoriaMetrics
|
|
// This function tries to do a simple query instead
|
|
log.Println("Warning: QueryDB called with Flux query, VictoriaMetrics uses MetricsQL")
|
|
log.Println("Query:", query)
|
|
|
|
// Try a simple health check instead
|
|
return nil, HealthCheck()
|
|
}
|