Compare commits

...

3 Commits

Author SHA1 Message Date
Max Kotliar
0ce9b3da1f app/vmselect: export traces to VictoriaTraces 2026-08-20 14:37:55 +03:00
Evgeny
a441d7e94e app/vmalert-tool: reuse connections to -remoteWrite.url
DebugClient sends every series in a separate request, but it left
MaxIdleConnsPerHost at the two connections of http.DefaultTransport.
Under concurrent rule evaluation most requests could not find an idle
connection and had to dial a new one, leaving many sockets in TIME_WAIT
state.

Set MaxIdleConnsPerHost from the new -remoteWrite.maxIdleConnections
flag, mirroring -datasource.maxIdleConnections, and apply the already
existing -remoteWrite.idleConnTimeout to the transport.
In my tests 640 concurrent pushes now open 17 connections instead of
227.

Related PR https://github.com/VictoriaMetrics/VictoriaMetrics/pull/11387/
2026-08-19 13:38:04 +02:00
Hui Wang
ac1d77e3de docs/changelog: add missing update note on v1.129.0 (#11430)
The breaking change was introduced with https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9779.
2026-08-19 12:05:23 +03:00
11 changed files with 581 additions and 1 deletions

View File

@@ -22,6 +22,7 @@ import (
"github.com/VictoriaMetrics/VictoriaMetrics/lib/procutil"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/promscrape"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/pushmetrics"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/querytracer/push"
)
var (
@@ -103,8 +104,10 @@ func main() {
logger.Infof("started VictoriaMetrics in %.3f seconds", time.Since(startTime).Seconds())
pushmetrics.Init()
push.Init()
sig := procutil.WaitForSigterm()
logger.Infof("received signal %s", sig)
push.Stop()
pushmetrics.Stop()
stopSelfScraper()

View File

@@ -36,6 +36,13 @@ func NewDebugClient() (*DebugClient, error) {
if err != nil {
return nil, fmt.Errorf("failed to create transport for -remoteWrite.url=%q: %w", *addr, err)
}
tr.IdleConnTimeout = *idleConnectionTimeout
// DebugClient sends every series in a separate request, so it needs more idle
// connections than the two http.DefaultTransport keeps per host.
tr.MaxIdleConnsPerHost = *maxIdleConnections
if tr.MaxIdleConns != 0 && tr.MaxIdleConns < tr.MaxIdleConnsPerHost {
tr.MaxIdleConns = tr.MaxIdleConnsPerHost
}
c := &DebugClient{
c: &http.Client{
Timeout: *sendTimeout,

View File

@@ -0,0 +1,45 @@
package remotewrite
import (
"net/http"
"testing"
)
// TestDebugClient_IdleConns makes sure DebugClient keeps enough idle connections
// to -remoteWrite.url. Every series is pushed in a separate request, so with the
// two idle connections per host of http.DefaultTransport most of the concurrent
// requests would open a new connection and leave a socket in TIME_WAIT state.
func TestDebugClient_IdleConns(t *testing.T) {
f := func(maxIdle int) {
t.Helper()
oldAddr, oldMaxIdle := *addr, *maxIdleConnections
*addr, *maxIdleConnections = "http://localhost:8428", maxIdle
defer func() {
*addr, *maxIdleConnections = oldAddr, oldMaxIdle
}()
client, err := NewDebugClient()
if err != nil {
t.Fatalf("failed to create debug client: %s", err)
}
tr, ok := client.c.Transport.(*http.Transport)
if !ok {
t.Fatalf("unexpected transport type %T", client.c.Transport)
}
if tr.MaxIdleConnsPerHost != maxIdle {
t.Fatalf("unexpected MaxIdleConnsPerHost; got %d; want %d", tr.MaxIdleConnsPerHost, maxIdle)
}
if tr.MaxIdleConns != 0 && tr.MaxIdleConns < maxIdle {
t.Fatalf("MaxIdleConns=%d is lower than MaxIdleConnsPerHost=%d", tr.MaxIdleConns, maxIdle)
}
if tr.IdleConnTimeout != *idleConnectionTimeout {
t.Fatalf("unexpected IdleConnTimeout; got %s; want %s", tr.IdleConnTimeout, *idleConnectionTimeout)
}
}
f(100)
// the number of idle connections must be raised together with the total limit
f(1000)
}

View File

@@ -34,6 +34,7 @@ var (
bearerTokenFile = flag.String("remoteWrite.bearerTokenFile", "", "Optional path to bearer token file to use for -remoteWrite.url.")
idleConnectionTimeout = flag.Duration("remoteWrite.idleConnTimeout", 50*time.Second, `Defines a duration for idle (keep-alive connections) to exist. Consider settings this value less to the value of "-http.idleConnTimeout". It must prevent possible "write: broken pipe" and "read: connection reset by peer" errors.`)
maxIdleConnections = flag.Int("remoteWrite.maxIdleConnections", 100, `Defines the number of idle (keep-alive connections) to -remoteWrite.url for the vmalert-tool debug writer, which sends every series in a separate request. Too low a value may result in a high number of sockets in TIME_WAIT state.`)
maxQueueSize = flag.Int("remoteWrite.maxQueueSize", defaultMaxQueueSize, "Defines the max number of pending datapoints to remote write endpoint")
maxBatchSize = flag.Int("remoteWrite.maxBatchSize", defaultMaxBatchSize, "Defines max number of timeseries to be flushed at once")

View File

@@ -26,6 +26,7 @@ import (
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/promscrape"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/querytracer"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/querytracer/push"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/timerpool"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/vmalertproxy"
)
@@ -117,8 +118,15 @@ func RequestHandler(w http.ResponseWriter, r *http.Request) bool {
// Handle non-trivial dynamic requests, which may take big amounts of time and resources.
startTime := time.Now()
defer requestDuration.UpdateDuration(startTime)
tracerEnabled := httputil.GetBool(r, "trace")
tracerEnabled := httputil.GetBool(r, "trace") || push.IsEnabled()
qt := querytracer.New(tracerEnabled, "%s", r.URL.Path)
if push.IsEnabled() {
defer func() {
if time.Since(startTime) >= push.MinTraceDuration() {
push.Push(qt)
}
}()
}
// Limit the number of concurrent queries.
select {

View File

@@ -29,6 +29,7 @@ See also [LTS releases](https://docs.victoriametrics.com/victoriametrics/lts-rel
* FEATURE: [vmui](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/#vmui): show the selected time zone UTC offset next to the date/time controls and allow opening time zone settings from it. See [#11332](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/11332).
* BUGFIX: [vmagent](https://docs.victoriametrics.com/victoriametrics/vmagent/) and `vminsert` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): fix infinite loop in the OpenTelemetry Firehose ingestion endpoint (`/opentelemetry/api/v1/push`) when receiving a malformed record with an incomplete varint in the `data` field. Previously this caused the goroutine to spin forever, permanently consuming CPU until the process was restarted.
* BUGFIX: [vmalert-tool](https://docs.victoriametrics.com/victoriametrics/vmalert-tool/): reuse connections to `-remoteWrite.url` when writing the results of recording rules and alerts. Previously every series was sent over a new connection, which left a lot of sockets in `TIME_WAIT` state and could exhaust the ephemeral port range. The number of idle connections can be tuned via the new `-remoteWrite.maxIdleConnections` command-line flag. Thanks @evkuzin for contribution.
* BUGFIX: [vmsingle](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/) and `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): prevent process crash in `sort_by_label_numeric()` and `sort_by_label_numeric_desc()` when a label value contains a number with 309 or more digits. See [#11423](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/11423).
## [v1.150.0](https://github.com/VictoriaMetrics/VictoriaMetrics/releases/tag/v1.150.0)

View File

@@ -115,6 +115,8 @@ Released at 2025-11-04
Released at 2025-10-31
**Update Note 1:** [vmalert](https://docs.victoriametrics.com/victoriametrics/vmalert/): reject responses with the [matrix](https://prometheus.io/docs/prometheus/latest/querying/basics/#expression-language-data-types) data type during normal rule evaluation, since vmalert expects the result to contain only a single sample or floating-point value as the rule value, not a matrix, which can contain a range of data points. Such responses could be generated by incorrect rule expressions such as `max_over_time(some_metric_filter > 90)[10m:]`, where `[10m:]` should be passed to `max_over_time` instead as `max_over_time((some_metric_filter > 90)[10m:])`.
* FEATURE: `vminsert` and `vmstorage` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): introduce new RPC protocol for insert-storage communication. See this PR [#9820](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/9820) for details.
* FEATURE: [vmalert](https://docs.victoriametrics.com/victoriametrics/vmalert/): explicitly check response type for [range queries](https://docs.victoriametrics.com/keyConcepts.html#range-query) during [replay](https://docs.victoriametrics.com/victoriametrics/vmalert/#rules-backfilling) and return error on type mismatch. This change should reduce confusions like in [#9779](https://github.com/VictoriaMetrics/VictoriaMetrics/issues/9779).
* FEATURE: [vmctl](https://docs.victoriametrics.com/victoriametrics/vmctl/): allow providing multiple filters for [remote-read migration mode](https://docs.victoriametrics.com/victoriametrics/vmctl/remoteread/) via multiple `--remote-read-filter-label` and `--remote-read-filter-label-value` flags. This is useful in order to narrow down the data being migrated by using more precise filters. See this PR [#9917](https://github.com/VictoriaMetrics/VictoriaMetrics/pull/9917) for details.

View File

@@ -370,6 +370,8 @@ See the docs at https://docs.victoriametrics.com/victoriametrics/vmalert/ .
Defines a duration for idle (keep-alive connections) to exist. Consider settings this value less to the value of "-http.idleConnTimeout". It must prevent possible "write: broken pipe" and "read: connection reset by peer" errors. (default 50s)
-remoteWrite.maxBatchSize int
Defines max number of timeseries to be flushed at once (default 10000)
-remoteWrite.maxIdleConnections int
Defines the number of idle (keep-alive connections) to -remoteWrite.url for the vmalert-tool debug writer, which sends every series in a separate request. Too low a value may result in a high number of sockets in TIME_WAIT state. (default 100)
-remoteWrite.maxQueueSize int
Defines the max number of pending datapoints to remote write endpoint (default 100000)
-remoteWrite.oauth2.clientID string

122
lib/querytracer/otlp.go Normal file
View File

@@ -0,0 +1,122 @@
package querytracer
import (
"crypto/rand"
"encoding/hex"
"time"
)
// OTLPTrace is an intermediate representation of a finished tracer tree
// suitable for conversion to OTLP protobuf spans.
type OTLPTrace struct {
TraceID string
Spans []OTLPSpan
}
// OTLPSpan represents a single span ready for OTLP export.
type OTLPSpan struct {
TraceID string
SpanID string
ParentSpanID string // empty for root
Name string
StartNano uint64
EndNano uint64
Events []OTLPEvent
}
// OTLPEvent is a timestamped annotation within a span (from Printf leaf nodes).
type OTLPEvent struct {
TimeNano uint64
Name string
}
// ToOTLPTrace converts the finished tracer tree into a flat list of OTLP spans.
// It must be called after Done/Donef.
func (t *Tracer) ToOTLPTrace() *OTLPTrace {
if t == nil {
return nil
}
traceID := newTraceID()
tr := &OTLPTrace{TraceID: traceID}
collectSpans(tr, t, traceID, "", t.startTime)
return tr
}
// collectSpans recursively walks the tracer tree and populates tr.Spans.
// parentSpanID is empty for the root.
// prevTime is used only when approximating timestamps for JSON-embedded spans.
func collectSpans(tr *OTLPTrace, t *Tracer, traceID, parentSpanID string, prevTime time.Time) {
if t.span != nil {
collectJSONSpans(tr, t.span, traceID, parentSpanID, prevTime)
return
}
isLeaf := t.doneTime.Equal(t.startTime)
if isLeaf {
// Printf leaf: no span to emit here; the caller handles attaching it.
return
}
// Regular span.
spanID := newSpanID()
spanIdx := len(tr.Spans)
tr.Spans = append(tr.Spans, OTLPSpan{
TraceID: traceID,
SpanID: spanID,
ParentSpanID: parentSpanID,
Name: t.message,
StartNano: uint64(t.startTime.UnixNano()),
EndNano: uint64(t.doneTime.UnixNano()),
})
// Process children. Leaf children become events on this span; non-leaf
// children recurse and produce their own spans.
childPrev := t.startTime
for _, child := range t.children {
if child.span == nil && child.doneTime.Equal(child.startTime) {
// Printf leaf: attach as span event.
tr.Spans[spanIdx].Events = append(tr.Spans[spanIdx].Events, OTLPEvent{
TimeNano: uint64(child.startTime.UnixNano()),
Name: child.message,
})
continue
}
collectSpans(tr, child, traceID, spanID, childPrev)
if !child.doneTime.IsZero() && !child.doneTime.Equal(child.startTime) {
childPrev = child.doneTime
}
}
}
// collectJSONSpans converts a span tree that was deserialized via AddJSON.
// Timestamps are approximated from prevTime since JSON spans only carry duration_msec.
func collectJSONSpans(tr *OTLPTrace, s *span, traceID, parentSpanID string, prevTime time.Time) {
startNano := uint64(prevTime.UnixNano())
durationNano := uint64(s.DurationMsec * float64(time.Millisecond))
spanID := newSpanID()
tr.Spans = append(tr.Spans, OTLPSpan{
TraceID: traceID,
SpanID: spanID,
ParentSpanID: parentSpanID,
Name: s.Message,
StartNano: startNano,
EndNano: startNano + durationNano,
})
childPrev := prevTime
for _, child := range s.Children {
collectJSONSpans(tr, child, traceID, spanID, childPrev)
childPrev = childPrev.Add(time.Duration(child.DurationMsec * float64(time.Millisecond)))
}
}
func newTraceID() string {
var b [16]byte
_, _ = rand.Read(b[:])
return hex.EncodeToString(b[:])
}
func newSpanID() string {
var b [8]byte
_, _ = rand.Read(b[:])
return hex.EncodeToString(b[:])
}

156
lib/querytracer/push/pb.go Normal file
View File

@@ -0,0 +1,156 @@
package push
// This file contains marshal-only protobuf types for OTLP trace export.
// Copied and trimmed from VictoriaTraces/lib/protoparser/opentelemetry/pb/
// Only marshaling is needed since this package only pushes traces outward.
import (
"encoding/hex"
"github.com/VictoriaMetrics/easyproto"
)
var mp easyproto.MarshalerPool
// exportTraceServiceRequest is the top-level OTLP protobuf message.
type exportTraceServiceRequest struct {
ResourceSpans []*resourceSpans
}
func (r *exportTraceServiceRequest) marshalProtobuf(dst []byte) []byte {
m := mp.Get()
mm := m.MessageMarshaler()
for _, rs := range r.ResourceSpans {
rs.marshalProtobuf(mm.AppendMessage(1))
}
dst = m.Marshal(dst)
mp.Put(m)
return dst
}
// resourceSpans groups spans from a single resource (e.g. service instance).
type resourceSpans struct {
Resource resource
ScopeSpans []*scopeSpans
}
func (rs *resourceSpans) marshalProtobuf(mm *easyproto.MessageMarshaler) {
rs.Resource.marshalProtobuf(mm.AppendMessage(1))
for _, ss := range rs.ScopeSpans {
ss.marshalProtobuf(mm.AppendMessage(2))
}
}
// resource holds resource-level attributes (service.name, etc.).
type resource struct {
Attributes []*keyValue
}
func (r *resource) marshalProtobuf(mm *easyproto.MessageMarshaler) {
for _, a := range r.Attributes {
a.marshalProtobuf(mm.AppendMessage(1))
}
}
// scopeSpans groups spans from a single instrumentation scope.
type scopeSpans struct {
Scope instrumentationScope
Spans []*span
}
func (ss *scopeSpans) marshalProtobuf(mm *easyproto.MessageMarshaler) {
ss.Scope.marshalProtobuf(mm.AppendMessage(1))
for _, s := range ss.Spans {
s.marshalProtobuf(mm.AppendMessage(2))
}
}
// instrumentationScope identifies the library that produced the spans.
type instrumentationScope struct {
Name string
Version string
}
func (is *instrumentationScope) marshalProtobuf(mm *easyproto.MessageMarshaler) {
mm.AppendString(1, is.Name)
mm.AppendString(2, is.Version)
}
// span represents a single operation within a trace.
type span struct {
// TraceID is a 32-char lowercase hex string (16 bytes).
TraceID string
// SpanID is a 16-char lowercase hex string (8 bytes).
SpanID string
// ParentSpanID is a 16-char lowercase hex string; empty for root spans.
ParentSpanID string
Name string
StartTimeUnixNano uint64
EndTimeUnixNano uint64
Attributes []*keyValue
Events []*spanEvent
}
func (s *span) marshalProtobuf(mm *easyproto.MessageMarshaler) {
traceID, err := hex.DecodeString(s.TraceID)
if err != nil {
traceID = []byte(s.TraceID)
}
mm.AppendBytes(1, traceID)
spanID, err := hex.DecodeString(s.SpanID)
if err != nil {
spanID = []byte(s.SpanID)
}
mm.AppendBytes(2, spanID)
// field 3: trace_state — omitted
parentSpanID, err := hex.DecodeString(s.ParentSpanID)
if err != nil {
parentSpanID = []byte(s.ParentSpanID)
}
mm.AppendBytes(4, parentSpanID)
mm.AppendString(5, s.Name)
// field 6: kind — omitted (INTERNAL=1 is default)
mm.AppendFixed64(7, s.StartTimeUnixNano)
mm.AppendFixed64(8, s.EndTimeUnixNano)
for _, a := range s.Attributes {
a.marshalProtobuf(mm.AppendMessage(9))
}
for _, e := range s.Events {
e.marshalProtobuf(mm.AppendMessage(11))
}
}
// spanEvent is a time-stamped annotation within a span.
type spanEvent struct {
TimeUnixNano uint64
Name string
}
func (se *spanEvent) marshalProtobuf(mm *easyproto.MessageMarshaler) {
mm.AppendFixed64(1, se.TimeUnixNano)
mm.AppendString(2, se.Name)
}
// keyValue is an OTLP attribute key-value pair.
type keyValue struct {
Key string
Value anyValue
}
func (kv *keyValue) marshalProtobuf(mm *easyproto.MessageMarshaler) {
mm.AppendString(1, kv.Key)
kv.Value.marshalProtobuf(mm.AppendMessage(2))
}
// anyValue holds a single string attribute value (sufficient for our use case).
type anyValue struct {
StringValue string
}
func (av *anyValue) marshalProtobuf(mm *easyproto.MessageMarshaler) {
mm.AppendString(1, av.StringValue)
}

View File

@@ -0,0 +1,233 @@
// Package push implements background exporting of VictoriaMetrics query traces
// to VictoriaTraces in OTLP protobuf format over HTTP.
package push
import (
"bytes"
"compress/gzip"
"flag"
"fmt"
"net/http"
"sync"
"time"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/buildinfo"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/logger"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/querytracer"
"github.com/VictoriaMetrics/metrics"
)
var (
pushURL = flag.String("search.traceExportURL", "", "If set, query traces are exported to this URL in OTLP protobuf format. "+
"For example, -search.traceExportURL=http://victoria-traces:4318/insert/opentelemetry/v1/traces . "+
"See https://docs.victoriametrics.com/victoriametrics/query-tracing/")
minTraceDuration = flag.Duration("search.traceExportMinDuration", 0, "Minimum query duration for exporting traces via -search.traceExportURL. "+
"Traces for faster queries are dropped. 0 means all traces are exported.")
)
var (
pushesTotal = metrics.NewCounter(`vm_trace_export_pushes_total`)
errorsTotal = metrics.NewCounter(`vm_trace_export_errors_total`)
droppedTotal = metrics.NewCounter(`vm_trace_export_dropped_total`)
)
// queueCap is the maximum number of pending traces waiting to be exported.
const queueCap = 1000
var (
queue chan *querytracer.OTLPTrace
stopCh chan struct{}
wg sync.WaitGroup
)
// IsEnabled returns true when trace export is configured.
func IsEnabled() bool {
return *pushURL != ""
}
// MinTraceDuration returns the configured minimum trace duration threshold.
func MinTraceDuration() time.Duration {
return *minTraceDuration
}
// Init starts the background export goroutine. Must be called after flag.Parse and logger.Init.
func Init() {
if !IsEnabled() {
return
}
queue = make(chan *querytracer.OTLPTrace, queueCap)
stopCh = make(chan struct{})
wg.Add(1)
go func() {
defer wg.Done()
runExporter()
}()
logger.Infof("started query trace exporter to %s", *pushURL)
}
// Stop drains pending traces and stops the background goroutine.
// Must be called during graceful shutdown before the process exits.
func Stop() {
if !IsEnabled() {
return
}
close(stopCh)
wg.Wait()
}
// Push enqueues t for async export. t must be a finished tracer.
// Traces are silently dropped if the queue is full.
func Push(t *querytracer.Tracer) {
if !t.Enabled() {
return
}
tr := t.ToOTLPTrace()
if tr == nil {
return
}
select {
case queue <- tr:
default:
droppedTotal.Inc()
}
}
// runExporter reads from queue and sends batches to the configured endpoint.
func runExporter() {
ticker := time.NewTicker(time.Second)
defer ticker.Stop()
var batch []*querytracer.OTLPTrace
for {
select {
case tr := <-queue:
batch = append(batch, tr)
if len(batch) >= 100 {
flushBatch(batch)
batch = batch[:0]
}
case <-ticker.C:
if len(batch) > 0 {
flushBatch(batch)
batch = batch[:0]
}
case <-stopCh:
// Drain remaining items.
for {
select {
case tr := <-queue:
batch = append(batch, tr)
default:
if len(batch) > 0 {
flushBatch(batch)
}
return
}
}
}
}
}
// flushBatch serializes and POSTs a batch of traces.
func flushBatch(batch []*querytracer.OTLPTrace) {
data, err := marshalBatch(batch)
if err != nil {
logger.Errorf("cannot marshal trace batch: %s", err)
errorsTotal.Inc()
return
}
if err := postData(data); err != nil {
logger.Warnf("cannot export query traces to %s: %s", *pushURL, err)
errorsTotal.Inc()
return
}
pushesTotal.Add(len(batch))
}
// marshalBatch packs all traces into a single ExportTraceServiceRequest and gzip-compresses it.
func marshalBatch(batch []*querytracer.OTLPTrace) ([]byte, error) {
svcVersion := buildinfo.Version
// Build resource attributes once (same for all spans in this process).
resAttrs := []*keyValue{
{Key: "service.name", Value: anyValue{StringValue: "victoriametrics"}},
{Key: "service.version", Value: anyValue{StringValue: svcVersion}},
}
var pbSpans []*span
for _, tr := range batch {
for i := range tr.Spans {
s := &tr.Spans[i]
pbSpan := &span{
TraceID: s.TraceID,
SpanID: s.SpanID,
ParentSpanID: s.ParentSpanID,
Name: s.Name,
StartTimeUnixNano: s.StartNano,
EndTimeUnixNano: s.EndNano,
}
for _, ev := range s.Events {
pbSpan.Events = append(pbSpan.Events, &spanEvent{
TimeUnixNano: ev.TimeNano,
Name: ev.Name,
})
}
pbSpans = append(pbSpans, pbSpan)
}
}
req := &exportTraceServiceRequest{
ResourceSpans: []*resourceSpans{
{
Resource: resource{Attributes: resAttrs},
ScopeSpans: []*scopeSpans{
{
Scope: instrumentationScope{
Name: "victoriametrics/querytracer",
Version: svcVersion,
},
Spans: pbSpans,
},
},
},
},
}
raw := req.marshalProtobuf(nil)
var buf bytes.Buffer
gz := gzip.NewWriter(&buf)
if _, err := gz.Write(raw); err != nil {
return nil, fmt.Errorf("cannot gzip trace data: %w", err)
}
if err := gz.Close(); err != nil {
return nil, fmt.Errorf("cannot close gzip writer: %w", err)
}
return buf.Bytes(), nil
}
var httpClient = &http.Client{
Timeout: 10 * time.Second,
}
// postData sends gzip-compressed protobuf data to the configured endpoint.
func postData(data []byte) error {
req, err := http.NewRequest(http.MethodPost, *pushURL, bytes.NewReader(data))
if err != nil {
return fmt.Errorf("cannot create HTTP request: %w", err)
}
req.Header.Set("Content-Type", "application/x-protobuf")
req.Header.Set("Content-Encoding", "gzip")
resp, err := httpClient.Do(req)
if err != nil {
return fmt.Errorf("HTTP request failed: %w", err)
}
defer resp.Body.Close()
if resp.StatusCode/100 != 2 {
return fmt.Errorf("unexpected HTTP response status %d", resp.StatusCode)
}
return nil
}