Compare commits

...

1 Commits

Author SHA1 Message Date
f41gh7
977c33caf2 app/vmselect: cancel query request if client disconnects
This commit adds context with cancel to the query requests.
Which allows to stop query execution earlier and release resources.
It's especially useful for long running queries - such as export of
federate requests. Which have large timeout.

 This commit only implements narrow sub-set of cancel:

* refactors query methods to accept and check context
* adds additional cancel checks to query execution functions to cancel
  data processing.

 It doesn't implement storage side query cancel. Which could be
implemented later.
2026-08-18 12:29:59 +02:00
17 changed files with 441 additions and 308 deletions

View File

@@ -1,6 +1,7 @@
package clusternative
import (
"context"
"flag"
"fmt"
"sync"
@@ -33,8 +34,10 @@ var (
)
// NewVMSelectServer starts new server at the given addr, which serves vmselect requests from netstorage.
func NewVMSelectServer(addr string) (*vmselectapi.Server, error) {
api := &vmstorageAPI{}
func NewVMSelectServer(ctx context.Context, addr string) (*vmselectapi.Server, error) {
api := &vmstorageAPI{
ctx: ctx,
}
limits := vmselectapi.Limits{
MaxConcurrentRequests: *maxConcurrentRequests,
MaxConcurrentRequestsFlagName: "clusternative.maxConcurrentRequests",
@@ -45,23 +48,25 @@ func NewVMSelectServer(addr string) (*vmselectapi.Server, error) {
}
// vmstorageAPI impelements vmselectapi.API
type vmstorageAPI struct{}
type vmstorageAPI struct {
ctx context.Context
}
func (api *vmstorageAPI) InitSearch(qt *querytracer.Tracer, sq *storage.SearchQuery, deadline uint64) (vmselectapi.BlockIterator, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
bi := newBlockIterator(qt, true, sq, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
bi := newBlockIterator(ctx, qt, true, sq)
return bi, nil
}
func (api *vmstorageAPI) Tenants(qt *querytracer.Tracer, tr storage.TimeRange, deadline uint64) ([]string, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
res, err := netstorage.Tenants(qt, tr, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
res, err := netstorage.Tenants(ctx, qt, tr)
return res, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) SearchMetricNames(qt *querytracer.Tracer, sq *storage.SearchQuery, deadline uint64) ([]string, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
metricNames, _, err := netstorage.SearchMetricNames(qt, true, sq, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
metricNames, _, err := netstorage.SearchMetricNames(ctx, qt, true, sq)
return metricNames, wrapClusterNativeError(err)
}
@@ -69,8 +74,8 @@ func (api *vmstorageAPI) LabelValues(qt *querytracer.Tracer, sq *storage.SearchQ
if maxLabelValues <= 0 || maxLabelValues > *maxTagValues {
maxLabelValues = *maxTagValues
}
dl := searchutil.DeadlineFromTimestamp(deadline)
labelValues, _, err := netstorage.LabelValues(qt, true, labelName, sq, maxLabelValues, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
labelValues, _, err := netstorage.LabelValues(ctx, qt, true, labelName, sq, maxLabelValues)
return labelValues, wrapClusterNativeError(err)
}
@@ -79,8 +84,8 @@ func (api *vmstorageAPI) TagValueSuffixes(qt *querytracer.Tracer, accountID, pro
if maxSuffixes <= 0 || maxSuffixes > *maxTagValueSuffixesPerSearch {
maxSuffixes = *maxTagValueSuffixesPerSearch
}
dl := searchutil.DeadlineFromTimestamp(deadline)
suffixes, _, err := netstorage.TagValueSuffixes(qt, accountID, projectID, true, tr, tagKey, tagValuePrefix, delimiter, maxSuffixes, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
suffixes, _, err := netstorage.TagValueSuffixes(ctx, qt, accountID, projectID, true, tr, tagKey, tagValuePrefix, delimiter, maxSuffixes)
return suffixes, wrapClusterNativeError(err)
}
@@ -88,48 +93,48 @@ func (api *vmstorageAPI) LabelNames(qt *querytracer.Tracer, sq *storage.SearchQu
if maxLabelNames <= 0 || maxLabelNames > *maxTagKeys {
maxLabelNames = *maxTagKeys
}
dl := searchutil.DeadlineFromTimestamp(deadline)
labelNames, _, err := netstorage.LabelNames(qt, true, sq, maxLabelNames, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
labelNames, _, err := netstorage.LabelNames(ctx, qt, true, sq, maxLabelNames)
return labelNames, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) SeriesCount(qt *querytracer.Tracer, accountID, projectID uint32, deadline uint64) (uint64, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
seriesCount, _, err := netstorage.SeriesCount(qt, accountID, projectID, true, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
seriesCount, _, err := netstorage.SeriesCount(ctx, qt, accountID, projectID, true)
return seriesCount, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) TSDBStatus(qt *querytracer.Tracer, sq *storage.SearchQuery, focusLabel string, topN int, deadline uint64) (*storage.TSDBStatus, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
tsdbStatus, _, err := netstorage.TSDBStatus(qt, true, sq, focusLabel, topN, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
tsdbStatus, _, err := netstorage.TSDBStatus(ctx, qt, true, sq, focusLabel, topN)
return tsdbStatus, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) DeleteSeries(qt *querytracer.Tracer, sq *storage.SearchQuery, deadline uint64) (int, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
deletedTotal, err := netstorage.DeleteSeries(qt, sq, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
deletedTotal, err := netstorage.DeleteSeries(ctx, qt, sq)
return deletedTotal, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) RegisterMetricNames(qt *querytracer.Tracer, mrs []storage.MetricRow, deadline uint64) error {
dl := searchutil.DeadlineFromTimestamp(deadline)
return wrapClusterNativeError(netstorage.RegisterMetricNames(qt, mrs, dl))
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
return wrapClusterNativeError(netstorage.RegisterMetricNames(ctx, qt, mrs))
}
func (api *vmstorageAPI) ResetMetricNamesUsageStats(qt *querytracer.Tracer, deadline uint64) error {
dl := searchutil.DeadlineFromTimestamp(deadline)
return wrapClusterNativeError(netstorage.ResetMetricNamesStats(qt, dl))
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
return wrapClusterNativeError(netstorage.ResetMetricNamesStats(ctx, qt))
}
func (api *vmstorageAPI) GetMetricNamesUsageStats(qt *querytracer.Tracer, tt *storage.TenantToken, le, limit int, matchPattern string, deadline uint64) (metricnamestats.StatsResult, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
statResult, err := netstorage.GetMetricNamesStats(qt, tt, le, limit, matchPattern, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
statResult, err := netstorage.GetMetricNamesStats(ctx, qt, tt, le, limit, matchPattern)
return statResult, wrapClusterNativeError(err)
}
func (api *vmstorageAPI) GetMetadataRecords(qt *querytracer.Tracer, tt *storage.TenantToken, limit int, metricName string, deadline uint64) ([]*metricsmetadata.Row, error) {
dl := searchutil.DeadlineFromTimestamp(deadline)
meta, _, err := netstorage.GetMetricsMetadata(qt, tt, true, limit, metricName, dl)
ctx := searchutil.NewContextWithDeadlineTimestamp(api.ctx, deadline)
meta, _, err := netstorage.GetMetricsMetadata(ctx, qt, tt, true, limit, metricName)
return meta, wrapClusterNativeError(err)
}
@@ -146,9 +151,9 @@ type workItem struct {
doneCh chan struct{}
}
func newBlockIterator(qt *querytracer.Tracer, denyPartialResponse bool, sq *storage.SearchQuery, deadline searchutil.Deadline) *blockIterator {
func newBlockIterator(ctx searchutil.Context, qt *querytracer.Tracer, denyPartialResponse bool, sq *storage.SearchQuery) *blockIterator {
bi := getBlockIterator()
workers, processBlocks := netstorage.PrepareProcessRawBlocks(qt, denyPartialResponse, sq, deadline)
workers, processBlocks := netstorage.PrepareProcessRawBlocks(ctx, qt, denyPartialResponse, sq)
bi.workCh = make(chan workItem, workers)
bi.wis = slicesutil.SetLength(bi.wis, workers)
for i := range bi.wis {

View File

@@ -23,12 +23,12 @@ var maxGraphitePathExpressionLen = flag.Int("search.maxGraphitePathExpressionLen
"Longer expressions are truncated to prevent memory exhaustion on complex nested queries. Set to 0 to disable truncation.")
type evalConfig struct {
ctx searchutil.Context
at *auth.Token
startTime int64
endTime int64
storageStep int64
denyPartialResponse bool
deadline searchutil.Deadline
currentTime time.Time
@@ -181,14 +181,14 @@ func evalMetricExpr(ec *evalConfig, me *graphiteql.MetricExpr) (nextSeriesFunc,
}
func newNextSeriesForSearchQuery(ec *evalConfig, sq *storage.SearchQuery, expr graphiteql.Expr) (nextSeriesFunc, error) {
rss, _, err := netstorage.ProcessSearchQuery(nil, ec.denyPartialResponse, sq, ec.deadline)
rss, _, err := netstorage.ProcessSearchQuery(ec.ctx, nil, ec.denyPartialResponse, sq)
if err != nil {
return nil, fmt.Errorf("cannot fetch data for %q: %w", sq, err)
}
seriesCh := make(chan *series, cgroup.AvailableCPUs())
errCh := make(chan error, 1)
go func() {
err := rss.RunParallel(nil, func(rs *netstorage.Result, _ uint) error {
err := rss.RunParallel(ec.ctx, nil, func(rs *netstorage.Result, _ uint) error {
nameWithTags := getCanonicalPath(&rs.MetricName)
tags := unmarshalTags(nameWithTags)
s := &series{
@@ -201,8 +201,9 @@ func newNextSeriesForSearchQuery(ec *evalConfig, sq *storage.SearchQuery, expr g
}
s.summarize(aggrAvg, ec.startTime, ec.endTime, ec.storageStep, 0)
deadline := ec.ctx.Deadline()
// A negative or zero duration will cause timer.C to return immediately
remainingTimeout := ec.deadline.Deadline() - fasttime.UnixTimestamp()
remainingTimeout := deadline.Deadline() - fasttime.UnixTimestamp()
t := timerpool.Get(time.Duration(remainingTimeout) * time.Second)
defer timerpool.Put(t)

View File

@@ -1,6 +1,7 @@
package graphite
import (
"context"
"fmt"
"math"
"reflect"
@@ -9,11 +10,13 @@ import (
"time"
"github.com/VictoriaMetrics/VictoriaMetrics/app/vmselect/graphiteql"
"github.com/VictoriaMetrics/VictoriaMetrics/app/vmselect/searchutil"
"github.com/VictoriaMetrics/VictoriaMetrics/lib/auth"
)
func TestExecExprSuccess(t *testing.T) {
ec := &evalConfig{
ctx: searchutil.NewContext(context.Background(), searchutil.NewDeadline(time.Now(), time.Minute, "")),
at: &auth.Token{},
startTime: 120e3,
endTime: 210e3,

View File

@@ -28,7 +28,7 @@ var maxTagValueSuffixes = flag.Int("search.maxTagValueSuffixesPerSearch", 100e3,
//
// See https://graphite-api.readthedocs.io/en/latest/api.html#metrics-find
func MetricsFindHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
format := r.FormValue("format")
if format == "" {
format = "treejson"
@@ -81,7 +81,7 @@ func MetricsFindHandler(startTime time.Time, at *auth.Token, w http.ResponseWrit
MaxTimestamp: until,
}
denyPartialResponse := httputil.GetDenyPartialResponse(r)
paths, isPartial, err := metricsFind(at, denyPartialResponse, tr, label, "", query, delimiter[0], false, deadline)
paths, isPartial, err := metricsFind(ctx, at, denyPartialResponse, tr, label, "", query, delimiter[0], false)
if err != nil {
return err
}
@@ -123,7 +123,7 @@ func deduplicatePaths(paths []string) []string {
//
// See https://graphite-api.readthedocs.io/en/latest/api.html#metrics-expand
func MetricsExpandHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
queries := r.Form["query"]
if len(queries) == 0 {
return fmt.Errorf("missing `query` arg")
@@ -159,7 +159,7 @@ func MetricsExpandHandler(startTime time.Time, at *auth.Token, w http.ResponseWr
isPartialResponse := false
denyPartialResponse := httputil.GetDenyPartialResponse(r)
for _, query := range queries {
paths, isPartial, err := metricsFind(at, denyPartialResponse, tr, label, "", query, delimiter[0], true, deadline)
paths, isPartial, err := metricsFind(ctx, at, denyPartialResponse, tr, label, "", query, delimiter[0], true)
if err != nil {
return err
}
@@ -208,11 +208,11 @@ func MetricsExpandHandler(startTime time.Time, at *auth.Token, w http.ResponseWr
//
// See https://graphite-api.readthedocs.io/en/latest/api.html#metrics-index-json
func MetricsIndexHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
jsonp := r.FormValue("jsonp")
denyPartialResponse := httputil.GetDenyPartialResponse(r)
sq := storage.NewSearchQuery(at.AccountID, at.ProjectID, 0, math.MaxInt64, nil, 0)
metricNames, isPartial, err := netstorage.LabelValues(nil, denyPartialResponse, "__name__", sq, 0, deadline)
metricNames, isPartial, err := netstorage.LabelValues(ctx, nil, denyPartialResponse, "__name__", sq, 0)
if err != nil {
return fmt.Errorf(`cannot obtain metric names: %w`, err)
}
@@ -229,12 +229,12 @@ func MetricsIndexHandler(startTime time.Time, at *auth.Token, w http.ResponseWri
}
// metricsFind searches for label values that match the given qHead and qTail.
func metricsFind(at *auth.Token, denyPartialResponse bool, tr storage.TimeRange, label, qHead, qTail string, delimiter byte,
isExpand bool, deadline searchutil.Deadline) ([]string, bool, error) {
func metricsFind(ctx searchutil.Context, at *auth.Token, denyPartialResponse bool, tr storage.TimeRange, label, qHead, qTail string, delimiter byte,
isExpand bool) ([]string, bool, error) {
n := strings.IndexAny(qTail, "*{[")
if n < 0 {
query := qHead + qTail
suffixes, isPartial, err := netstorage.TagValueSuffixes(nil, at.AccountID, at.ProjectID, denyPartialResponse, tr, label, query, delimiter, *maxTagValueSuffixes, deadline)
suffixes, isPartial, err := netstorage.TagValueSuffixes(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, tr, label, query, delimiter, *maxTagValueSuffixes)
if err != nil {
return nil, false, err
}
@@ -254,7 +254,7 @@ func metricsFind(at *auth.Token, denyPartialResponse bool, tr storage.TimeRange,
}
if n == len(qTail)-1 && strings.HasSuffix(qTail, "*") {
query := qHead + qTail[:len(qTail)-1]
suffixes, isPartial, err := netstorage.TagValueSuffixes(nil, at.AccountID, at.ProjectID, denyPartialResponse, tr, label, query, delimiter, *maxTagValueSuffixes, deadline)
suffixes, isPartial, err := netstorage.TagValueSuffixes(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, tr, label, query, delimiter, *maxTagValueSuffixes)
if err != nil {
return nil, false, err
}
@@ -268,7 +268,7 @@ func metricsFind(at *auth.Token, denyPartialResponse bool, tr storage.TimeRange,
return results, isPartial, nil
}
qHead += qTail[:n]
paths, isPartial, err := metricsFind(at, denyPartialResponse, tr, label, qHead, "*", delimiter, isExpand, deadline)
paths, isPartial, err := metricsFind(ctx, at, denyPartialResponse, tr, label, qHead, "*", delimiter, isExpand)
if err != nil {
return nil, false, err
}
@@ -292,7 +292,7 @@ func metricsFind(at *auth.Token, denyPartialResponse bool, tr storage.TimeRange,
results = append(results, path)
continue
}
fullPaths, isPartialLocal, err := metricsFind(at, denyPartialResponse, tr, label, path, qTail, delimiter, isExpand, deadline)
fullPaths, isPartialLocal, err := metricsFind(ctx, at, denyPartialResponse, tr, label, path, qTail, delimiter, isExpand)
if err != nil {
return nil, false, err
}

View File

@@ -27,7 +27,7 @@ var (
//
// See https://graphite.readthedocs.io/en/stable/render_api.html
func RenderHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
format := r.FormValue("format")
if format != "json" {
return fmt.Errorf("unsupported format=%q; supported values: json", format)
@@ -100,12 +100,12 @@ func RenderHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r
targets := r.Form["target"]
for _, target := range targets {
ec := &evalConfig{
ctx: ctx,
at: at,
startTime: fromTime,
endTime: untilTime,
storageStep: storageStep,
denyPartialResponse: denyPartialResponse,
deadline: deadline,
currentTime: startTime,
xFilesFactor: xFilesFactor,
etfs: etfs,

View File

@@ -31,7 +31,7 @@ var (
//
// See https://graphite.readthedocs.io/en/stable/tags.html#removing-series-from-the-tagdb
func TagsDelSeriesHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
paths := r.Form["path"]
totalDeleted := 0
var row graphiteparser.Row
@@ -60,7 +60,7 @@ func TagsDelSeriesHandler(startTime time.Time, at *auth.Token, w http.ResponseWr
}
tfss := joinTagFilterss(tfs, etfs)
sq := storage.NewSearchQuery(at.AccountID, at.ProjectID, 0, ct, tfss, 0)
n, err := netstorage.DeleteSeries(nil, sq, deadline)
n, err := netstorage.DeleteSeries(ctx, nil, sq)
if err != nil {
return fmt.Errorf("cannot delete series for %q: %w", sq, err)
}
@@ -91,7 +91,7 @@ func TagsTagMultiSeriesHandler(startTime time.Time, at *auth.Token, w http.Respo
}
func registerMetrics(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request, isJSONResponse bool) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
paths := r.Form["path"]
var row graphiteparser.Row
var labels []prompb.Label
@@ -137,7 +137,7 @@ func registerMetrics(startTime time.Time, at *auth.Token, w http.ResponseWriter,
mr.MetricNameRaw = storage.MarshalMetricNameRaw(mr.MetricNameRaw[:0], at.AccountID, at.ProjectID, labels)
mr.Timestamp = ct
}
if err := netstorage.RegisterMetricNames(nil, mrs, deadline); err != nil {
if err := netstorage.RegisterMetricNames(ctx, nil, mrs); err != nil {
return fmt.Errorf("cannot register paths: %w", err)
}
@@ -165,7 +165,7 @@ var (
//
// See https://graphite.readthedocs.io/en/stable/tags.html#auto-complete-support
func TagsAutoCompleteValuesHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
limit, err := httputil.GetInt(r, "limit")
if err != nil {
return err
@@ -192,7 +192,7 @@ func TagsAutoCompleteValuesHandler(startTime time.Time, at *auth.Token, w http.R
// Escape special chars in tagPrefix as Graphite does.
// See https://github.com/graphite-project/graphite-web/blob/3ad279df5cb90b211953e39161df416e54a84948/webapp/graphite/tags/base.py#L228
filter := regexp.QuoteMeta(valuePrefix)
tagValues, isPartial, err = netstorage.GraphiteTagValues(nil, at.AccountID, at.ProjectID, denyPartialResponse, tag, filter, *maxGraphiteTagValuesPerSearch, deadline)
tagValues, isPartial, err = netstorage.GraphiteTagValues(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, tag, filter, *maxGraphiteTagValuesPerSearch)
if err != nil {
return err
}
@@ -202,7 +202,7 @@ func TagsAutoCompleteValuesHandler(startTime time.Time, at *auth.Token, w http.R
if err != nil {
return err
}
metricNames, isPartialResponse, err := netstorage.SearchMetricNames(nil, denyPartialResponse, sq, deadline)
metricNames, isPartialResponse, err := netstorage.SearchMetricNames(ctx, nil, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch metric names for %q: %w", sq, err)
}
@@ -258,7 +258,7 @@ var tagsAutoCompleteValuesDuration = metrics.NewSummary(`vm_request_duration_sec
//
// See https://graphite.readthedocs.io/en/stable/tags.html#auto-complete-support
func TagsAutoCompleteTagsHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
limit, err := httputil.GetInt(r, "limit")
if err != nil {
return err
@@ -282,7 +282,7 @@ func TagsAutoCompleteTagsHandler(startTime time.Time, at *auth.Token, w http.Res
// Escape special chars in tagPrefix as Graphite does.
// See https://github.com/graphite-project/graphite-web/blob/3ad279df5cb90b211953e39161df416e54a84948/webapp/graphite/tags/base.py#L181
filter := regexp.QuoteMeta(tagPrefix)
labels, isPartial, err = netstorage.GraphiteTags(nil, at.AccountID, at.ProjectID, denyPartialResponse, filter, *maxGraphiteTagKeysPerSearch, deadline)
labels, isPartial, err = netstorage.GraphiteTags(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, filter, *maxGraphiteTagKeysPerSearch)
if err != nil {
return err
}
@@ -292,7 +292,7 @@ func TagsAutoCompleteTagsHandler(startTime time.Time, at *auth.Token, w http.Res
if err != nil {
return err
}
metricNames, isPartialResponse, err := netstorage.SearchMetricNames(nil, denyPartialResponse, sq, deadline)
metricNames, isPartialResponse, err := netstorage.SearchMetricNames(ctx, nil, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch metric names for %q: %w", sq, err)
}
@@ -344,7 +344,7 @@ var tagsAutoCompleteTagsDuration = metrics.NewSummary(`vm_request_duration_secon
//
// See https://graphite.readthedocs.io/en/stable/tags.html#exploring-tags
func TagsFindSeriesHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
limit, err := httputil.GetInt(r, "limit")
if err != nil {
return err
@@ -362,7 +362,7 @@ func TagsFindSeriesHandler(startTime time.Time, at *auth.Token, w http.ResponseW
return err
}
denyPartialResponse := httputil.GetDenyPartialResponse(r)
metricNames, isPartial, err := netstorage.SearchMetricNames(nil, denyPartialResponse, sq, deadline)
metricNames, isPartial, err := netstorage.SearchMetricNames(ctx, nil, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch metric names for %q: %w", sq, err)
}
@@ -420,14 +420,14 @@ var tagsFindSeriesDuration = metrics.NewSummary(`vm_request_duration_seconds{pat
//
// See https://graphite.readthedocs.io/en/stable/tags.html#exploring-tags
func TagValuesHandler(startTime time.Time, at *auth.Token, tagName string, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
limit, err := httputil.GetInt(r, "limit")
if err != nil {
return err
}
filter := r.FormValue("filter")
denyPartialResponse := httputil.GetDenyPartialResponse(r)
tagValues, isPartial, err := netstorage.GraphiteTagValues(nil, at.AccountID, at.ProjectID, denyPartialResponse, tagName, filter, *maxGraphiteTagValuesPerSearch, deadline)
tagValues, isPartial, err := netstorage.GraphiteTagValues(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, tagName, filter, *maxGraphiteTagValuesPerSearch)
if err != nil {
return err
}
@@ -452,14 +452,14 @@ var tagValuesDuration = metrics.NewSummary(`vm_request_duration_seconds{path="/t
//
// See https://graphite.readthedocs.io/en/stable/tags.html#exploring-tags
func TagsHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
limit, err := httputil.GetInt(r, "limit")
if err != nil {
return err
}
filter := r.FormValue("filter")
denyPartialResponse := httputil.GetDenyPartialResponse(r)
labels, isPartial, err := netstorage.GraphiteTags(nil, at.AccountID, at.ProjectID, denyPartialResponse, filter, *maxGraphiteTagKeysPerSearch, deadline)
labels, isPartial, err := netstorage.GraphiteTags(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse, filter, *maxGraphiteTagKeysPerSearch)
if err != nil {
return err
}

View File

@@ -1,6 +1,7 @@
package main
import (
"context"
"embed"
"encoding/json"
"flag"
@@ -101,6 +102,8 @@ func main() {
buildinfo.Init()
logger.Init()
ctx, cancel := context.WithCancel(context.Background())
defer cancel()
logger.Infof("starting netstorage at storageNodes %s", *storageNodes)
startTime := time.Now()
storage.SetDedupInterval(*minScrapeInterval)
@@ -135,7 +138,7 @@ func main() {
var vmselectapiServer *vmselectapi.Server
if *clusternativeListenAddr != "" {
logger.Infof("starting vmselectapi server at %q", *clusternativeListenAddr)
s, err := clusternative.NewVMSelectServer(*clusternativeListenAddr)
s, err := clusternative.NewVMSelectServer(ctx, *clusternativeListenAddr)
if err != nil {
logger.Fatalf("cannot initialize vmselectapi server: %s", err)
}
@@ -162,6 +165,7 @@ func main() {
logger.Fatalf("cannot stop http service: %s", err)
}
logger.Infof("successfully shut down http service in %.3f seconds", time.Since(startTime).Seconds())
cancel()
if vmselectapiServer != nil {
logger.Infof("stopping vmselectapi server...")
@@ -196,7 +200,6 @@ var (
func requestHandler(w http.ResponseWriter, r *http.Request) bool {
path := strings.ReplaceAll(r.URL.Path, "//", "/")
if handleStaticAndSimpleRequests(w, r, path) {
return true
}

File diff suppressed because it is too large Load Diff

View File

@@ -21,7 +21,7 @@ var (
)
// TenantsCached returns the list of tenants available in the storage.
func TenantsCached(qt *querytracer.Tracer, tr storage.TimeRange, deadline searchutil.Deadline, mayCache bool) ([]storage.TenantToken, error) {
func TenantsCached(ctx searchutil.Context, qt *querytracer.Tracer, tr storage.TimeRange, mayCache bool) ([]storage.TenantToken, error) {
qtL := qt.NewChild("fetching tenants on timeRange=%s", tr.String())
defer qtL.Done()
@@ -41,7 +41,7 @@ func TenantsCached(qt *querytracer.Tracer, tr storage.TimeRange, deadline search
qtL.Printf("do not fetch list of tenants from cache")
}
tenants, err := Tenants(qtL, tr, deadline)
tenants, err := Tenants(ctx, qtL, tr)
if err != nil {
return nil, fmt.Errorf("cannot obtain tenants: %w", err)
}

View File

@@ -13,8 +13,8 @@ import (
)
// GetTenantTokensFromFilters returns the list of tenant tokens and the list of filters without tenant filters.
func GetTenantTokensFromFilters(qt *querytracer.Tracer, tr storage.TimeRange, tfs [][]storage.TagFilter, deadline searchutil.Deadline, mayCache bool) ([]storage.TenantToken, [][]storage.TagFilter, error) {
tenants, err := TenantsCached(qt, tr, deadline, mayCache)
func GetTenantTokensFromFilters(ctx searchutil.Context, qt *querytracer.Tracer, tr storage.TimeRange, tfs [][]storage.TagFilter, mayCache bool) ([]storage.TenantToken, [][]storage.TagFilter, error) {
tenants, err := TenantsCached(ctx, qt, tr, mayCache)
if err != nil {
return nil, nil, fmt.Errorf("cannot obtain tenants: %w", err)
}

View File

@@ -144,7 +144,7 @@ func FederateHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter,
return fmt.Errorf("cannot obtain search query: %w", err)
}
denyPartialResponse := httputil.GetDenyPartialResponse(r)
rss, isPartial, err := netstorage.ProcessSearchQuery(nil, denyPartialResponse, sq, cp.deadline)
rss, isPartial, err := netstorage.ProcessSearchQuery(cp.ctx, nil, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch data for %q: %w", sq, err)
}
@@ -172,7 +172,7 @@ func FederateHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter,
bw := bufferedwriter.Get(w)
defer bufferedwriter.Put(bw)
sw := newScalableWriter(bw)
err = rss.RunParallel(nil, func(rs *netstorage.Result, workerID uint) error {
err = rss.RunParallel(cp.ctx, nil, func(rs *netstorage.Result, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -229,12 +229,12 @@ func ExportCSVHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter
// Unconditionally deny partial response for the exported data,
// since users usually expect that the exported data is full.
denyPartialResponse := true
rss, _, err := netstorage.ProcessSearchQuery(nil, denyPartialResponse, sq, cp.deadline)
rss, _, err := netstorage.ProcessSearchQuery(cp.ctx, nil, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch data for %q: %w", sq, err)
}
go func() {
err := rss.RunParallel(nil, func(rs *netstorage.Result, workerID uint) error {
err := rss.RunParallel(cp.ctx, nil, func(rs *netstorage.Result, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -253,7 +253,7 @@ func ExportCSVHandler(startTime time.Time, at *auth.Token, w http.ResponseWriter
}()
} else {
go func() {
err := netstorage.ExportBlocks(nil, sq, cp.deadline, func(mn *storage.MetricName, b *storage.Block, tr storage.TimeRange, workerID uint) error {
err := netstorage.ExportBlocks(cp.ctx, nil, sq, func(mn *storage.MetricName, b *storage.Block, tr storage.TimeRange, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -310,7 +310,7 @@ func ExportNativeHandler(startTime time.Time, at *auth.Token, w http.ResponseWri
_, _ = bw.Write(trBuf)
// Marshal native blocks.
err = netstorage.ExportBlocks(nil, sq, cp.deadline, func(mn *storage.MetricName, b *storage.Block, _ storage.TimeRange, workerID uint) error {
err = netstorage.ExportBlocks(cp.ctx, nil, sq, func(mn *storage.MetricName, b *storage.Block, _ storage.TimeRange, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -455,13 +455,13 @@ func exportHandler(qt *querytracer.Tracer, at *auth.Token, w http.ResponseWriter
// Unconditionally deny partial response for the exported data,
// since users usually expect that the exported data is full.
denyPartialResponse := true
rss, _, err := netstorage.ProcessSearchQuery(qt, denyPartialResponse, sq, cp.deadline)
rss, _, err := netstorage.ProcessSearchQuery(cp.ctx, qt, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch data for %q: %w", sq, err)
}
qtChild := qt.NewChild("background export format=%s", format)
go func() {
err := rss.RunParallel(qtChild, func(rs *netstorage.Result, workerID uint) error {
err := rss.RunParallel(cp.ctx, qtChild, func(rs *netstorage.Result, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -482,7 +482,7 @@ func exportHandler(qt *querytracer.Tracer, at *auth.Token, w http.ResponseWriter
} else {
qtChild := qt.NewChild("background export format=%s", format)
go func() {
err := netstorage.ExportBlocks(qtChild, sq, cp.deadline, func(mn *storage.MetricName, b *storage.Block, tr storage.TimeRange, workerID uint) error {
err := netstorage.ExportBlocks(cp.ctx, qtChild, sq, func(mn *storage.MetricName, b *storage.Block, tr storage.TimeRange, workerID uint) error {
if err := bw.Error(); err != nil {
return err
}
@@ -549,7 +549,7 @@ func DeleteHandler(startTime time.Time, at *auth.Token, r *http.Request) error {
if err != nil {
return err
}
cp.deadline = searchutil.GetDeadlineForDelete(r, startTime)
cp.ctx = searchutil.GetContextForDelete(r, startTime)
if !cp.IsDefaultTimeRange() {
return fmt.Errorf("delete API does not support specific time ranges using start and end args, the series can only be deleted completely")
@@ -558,7 +558,7 @@ func DeleteHandler(startTime time.Time, at *auth.Token, r *http.Request) error {
if err != nil {
return err
}
deletedCount, err := netstorage.DeleteSeries(nil, sq, cp.deadline)
deletedCount, err := netstorage.DeleteSeries(cp.ctx, nil, sq)
if err != nil {
return fmt.Errorf("cannot delete time series: %w", err)
}
@@ -649,7 +649,7 @@ var httpClient = &http.Client{
// Tenants processes /admin/tenants request.
func Tenants(qt *querytracer.Tracer, startTime time.Time, w http.ResponseWriter, r *http.Request) error {
deadline := searchutil.GetDeadlineForStatusRequest(r, startTime)
ctx := searchutil.GetContextForStatusRequest(r, startTime)
start, err := httputil.GetTime(r, "start", 0)
if err != nil {
return err
@@ -663,7 +663,7 @@ func Tenants(qt *querytracer.Tracer, startTime time.Time, w http.ResponseWriter,
MinTimestamp: start,
MaxTimestamp: end,
}
tenants, err := netstorage.Tenants(qt, tr, deadline)
tenants, err := netstorage.Tenants(ctx, qt, tr)
if err != nil {
return err
}
@@ -702,7 +702,7 @@ func LabelValuesHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.To
// Spec: https://github.com/prometheus/proposals/blob/main/proposals/0028-utf8.md
labelName = unescapePrometheusLabelName(labelName)
}
labelValues, isPartial, err := netstorage.LabelValues(qt, denyPartialResponse, labelName, sq, limit, cp.deadline)
labelValues, isPartial, err := netstorage.LabelValues(cp.ctx, qt, denyPartialResponse, labelName, sq, limit)
if err != nil {
return fmt.Errorf("cannot obtain values for label %q: %w", labelName, err)
}
@@ -732,7 +732,7 @@ func TSDBStatusHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
if err != nil {
return httpserver.InvalidParamError(err)
}
cp.deadline = searchutil.GetDeadlineForStatusRequest(r, startTime)
cp.ctx = searchutil.GetContextForStatusRequest(r, startTime)
date := fasttime.UnixDate()
dateStr := r.FormValue("date")
@@ -770,7 +770,7 @@ func TSDBStatusHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
if err != nil {
return err
}
status, isPartial, err := netstorage.TSDBStatus(qt, denyPartialResponse, sq, focusLabel, topN, cp.deadline)
status, isPartial, err := netstorage.TSDBStatus(cp.ctx, qt, denyPartialResponse, sq, focusLabel, topN)
if err != nil {
return fmt.Errorf("cannot obtain tsdb stats: %w", err)
}
@@ -806,7 +806,7 @@ func LabelsHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token,
if err != nil {
return err
}
labels, isPartial, err := netstorage.LabelNames(qt, denyPartialResponse, sq, limit, cp.deadline)
labels, isPartial, err := netstorage.LabelNames(cp.ctx, qt, denyPartialResponse, sq, limit)
if err != nil {
return fmt.Errorf("cannot obtain labels: %w", err)
}
@@ -850,7 +850,7 @@ func MetadataHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token
denyPartialResponse := httputil.GetDenyPartialResponse(r)
metadata, isPartial, err := netstorage.GetMetricsMetadata(qt, tt, denyPartialResponse, limit, metricName, cp.deadline)
metadata, isPartial, err := netstorage.GetMetricsMetadata(cp.ctx, qt, tt, denyPartialResponse, limit, metricName)
if err != nil {
return fmt.Errorf("cannot get metadata: %w", err)
}
@@ -881,7 +881,7 @@ func getSearchQuery(qt *querytracer.Tracer, at *auth.Token, cp *commonParams, ma
if at != nil {
return storage.NewSearchQuery(at.AccountID, at.ProjectID, cp.start, cp.end, cp.filterss, maxSeries), nil
}
tt, tfs, err := netstorage.GetTenantTokensFromFilters(qt, storage.TimeRange{MinTimestamp: cp.start, MaxTimestamp: cp.end}, cp.filterss, cp.deadline, true)
tt, tfs, err := netstorage.GetTenantTokensFromFilters(cp.ctx, qt, storage.TimeRange{MinTimestamp: cp.start, MaxTimestamp: cp.end}, cp.filterss, true)
if err != nil {
return nil, fmt.Errorf("cannot obtain tenant tokens: %w", err)
}
@@ -897,9 +897,9 @@ func SeriesCountHandler(startTime time.Time, at *auth.Token, w http.ResponseWrit
if at == nil {
return fmt.Errorf("multi-tenant request to /api/v1/series/count is not supported")
}
deadline := searchutil.GetDeadlineForStatusRequest(r, startTime)
ctx := searchutil.GetContextForStatusRequest(r, startTime)
denyPartialResponse := httputil.GetDenyPartialResponse(r)
n, isPartial, err := netstorage.SeriesCount(nil, at.AccountID, at.ProjectID, denyPartialResponse, deadline)
n, isPartial, err := netstorage.SeriesCount(ctx, nil, at.AccountID, at.ProjectID, denyPartialResponse)
if err != nil {
return fmt.Errorf("cannot obtain series count: %w", err)
}
@@ -940,7 +940,7 @@ func SeriesHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token,
return err
}
denyPartialResponse := httputil.GetDenyPartialResponse(r)
metricNames, isPartial, err := netstorage.SearchMetricNames(qt, denyPartialResponse, sq, cp.deadline)
metricNames, isPartial, err := netstorage.SearchMetricNames(cp.ctx, qt, denyPartialResponse, sq)
if err != nil {
return fmt.Errorf("cannot fetch time series for %q: %w", sq, err)
}
@@ -966,7 +966,7 @@ func QueryHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token, w
defer queryDuration.UpdateDuration(startTime)
ct := startTime.UnixNano() / 1e6
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
noCache := httputil.GetBool(r, "nocache")
query := r.FormValue("query")
if len(query) == 0 {
@@ -1018,7 +1018,7 @@ func QueryHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token, w
filterss := searchutil.JoinTagFilterss(tagFilterss, etfs)
cp := &commonParams{
deadline: deadline,
ctx: ctx,
start: start,
end: end,
filterss: filterss,
@@ -1067,13 +1067,13 @@ func QueryHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token, w
queryOffset = 0
}
ec := &promql.EvalConfig{
Context: ctx,
Start: start,
End: start,
Step: step,
MaxPointsPerSeries: *maxPointsPerTimeseries,
MaxSeries: *maxUniqueTimeseries,
QuotedRemoteAddr: httpserver.GetQuotedRemoteAddr(r),
Deadline: deadline,
NoCache: noCache,
LookbackDelta: lookbackDelta,
RoundDigits: getRoundDigits(r),
@@ -1085,7 +1085,7 @@ func QueryHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token, w
DenyPartialResponse: httputil.GetDenyPartialResponse(r),
}
err = populateAuthTokens(qt, ec, at, deadline)
err = populateAuthTokens(ctx, qt, ec, at)
if err != nil {
return fmt.Errorf("cannot populate auth tokens: %w", err)
}
@@ -1169,7 +1169,7 @@ func QueryRangeHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
func queryRangeHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Token, w http.ResponseWriter, query string,
start, end, step, lookbackDelta int64, r *http.Request, ct int64, etfs [][]storage.TagFilter) error {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
noCache := httputil.GetBool(r, "nocache")
optimizeRepeatedBinaryOpSubexprs := httputil.GetBool(r, "optimize_repeated_binary_op_subexprs")
if start > end {
@@ -1183,13 +1183,13 @@ func queryRangeHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
}
ec := &promql.EvalConfig{
Context: ctx,
Start: start,
End: end,
Step: step,
MaxPointsPerSeries: *maxPointsPerTimeseries,
MaxSeries: *maxUniqueTimeseries,
QuotedRemoteAddr: httpserver.GetQuotedRemoteAddr(r),
Deadline: deadline,
NoCache: noCache,
OptimizeRepeatedBinaryOpSubexprs: optimizeRepeatedBinaryOpSubexprs,
LookbackDelta: lookbackDelta,
@@ -1202,7 +1202,8 @@ func queryRangeHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
DenyPartialResponse: httputil.GetDenyPartialResponse(r),
}
if err := populateAuthTokens(qt, ec, at, deadline); err != nil {
if err := populateAuthTokens(ctx, qt, ec, at); err != nil {
return fmt.Errorf("cannot populate auth tokens: %w", err)
}
qs := promql.NewQueryStats(query, at, ec)
@@ -1240,13 +1241,13 @@ func queryRangeHandler(qt *querytracer.Tracer, startTime time.Time, at *auth.Tok
return nil
}
func populateAuthTokens(qt *querytracer.Tracer, ec *promql.EvalConfig, at *auth.Token, deadline searchutil.Deadline) error {
func populateAuthTokens(ctx searchutil.Context, qt *querytracer.Tracer, ec *promql.EvalConfig, at *auth.Token) error {
if at != nil {
ec.AuthTokens = []*auth.Token{at}
return nil
}
tt, tfs, err := netstorage.GetTenantTokensFromFilters(qt, storage.TimeRange{MinTimestamp: ec.Start, MaxTimestamp: ec.End}, ec.EnforcedTagFilterss, deadline, ec.MayCache())
tt, tfs, err := netstorage.GetTenantTokensFromFilters(ctx, qt, storage.TimeRange{MinTimestamp: ec.Start, MaxTimestamp: ec.End}, ec.EnforcedTagFilterss, ec.MayCache())
if err != nil {
return fmt.Errorf("cannot obtain tenant tokens for the given search query: %w", err)
}
@@ -1417,7 +1418,7 @@ func QueryStatsHandler(at *auth.Token, w http.ResponseWriter, r *http.Request) e
//
// timeout, start, end, match[], extra_label, extra_filters[]
type commonParams struct {
deadline searchutil.Deadline
ctx searchutil.Context
start int64
end int64
currentTimestamp int64
@@ -1441,7 +1442,7 @@ func getExportParams(r *http.Request, startTime time.Time) (*commonParams, error
if err != nil {
return nil, err
}
cp.deadline = searchutil.GetDeadlineForExport(r, startTime)
cp.ctx = searchutil.GetContextForExport(r, startTime)
return cp, nil
}
@@ -1453,7 +1454,7 @@ func getCommonParamsForLabelsAPI(r *http.Request, startTime time.Time, requireNo
if cp.start == 0 {
cp.start = cp.end - defaultStep
}
cp.deadline = searchutil.GetDeadlineForLabelsAPI(r, startTime)
cp.ctx = searchutil.GetContextForLabelsAPI(r, startTime)
return cp, nil
}
@@ -1470,7 +1471,7 @@ func getCommonParams(r *http.Request, startTime time.Time, requireNonEmptyMatch
}
func getCommonParamsInternal(r *http.Request, startTime time.Time, requireNonEmptyMatch, isLabelsAPI bool) (*commonParams, error) {
deadline := searchutil.GetDeadlineForQuery(r, startTime)
ctx := searchutil.GetContextForQuery(r, startTime)
start, err := httputil.GetTime(r, "start", 0)
if err != nil {
return nil, err
@@ -1510,7 +1511,7 @@ func getCommonParamsInternal(r *http.Request, startTime time.Time, requireNonEmp
}
cp := &commonParams{
deadline: deadline,
ctx: ctx,
start: start,
end: end,
currentTimestamp: ct,

View File

@@ -114,7 +114,10 @@ func alignStartEnd(start, end, step int64) (int64, int64) {
// EvalConfig is the configuration required for query evaluation via Exec
type EvalConfig struct {
AuthTokens []*auth.Token
Context searchutil.Context
AuthTokens []*auth.Token
IsMultiTenant bool
Start int64
@@ -131,8 +134,6 @@ type EvalConfig struct {
// QuotedRemoteAddr contains quoted remote address.
QuotedRemoteAddr string
Deadline searchutil.Deadline
// Whether the response must not be cached.
NoCache bool
@@ -177,6 +178,7 @@ type EvalConfig struct {
// copyEvalConfig returns src copy.
func copyEvalConfig(src *EvalConfig) *EvalConfig {
var ec EvalConfig
ec.Context = src.Context
ec.AuthTokens = src.AuthTokens
ec.IsMultiTenant = src.IsMultiTenant
ec.Start = src.Start
@@ -184,7 +186,6 @@ func copyEvalConfig(src *EvalConfig) *EvalConfig {
ec.Step = src.Step
ec.MaxSeries = src.MaxSeries
ec.MaxPointsPerSeries = src.MaxPointsPerSeries
ec.Deadline = src.Deadline
ec.NoCache = src.NoCache
ec.OptimizeRepeatedBinaryOpSubexprs = src.OptimizeRepeatedBinaryOpSubexprs
ec.LookbackDelta = src.LookbackDelta
@@ -1877,7 +1878,7 @@ func evalRollupFuncNoCache(qt *querytracer.Tracer, ec *EvalConfig, funcName stri
} else {
sq = storage.NewSearchQuery(ec.AuthTokens[0].AccountID, ec.AuthTokens[0].ProjectID, minTimestamp, ec.End, tfss, ec.MaxSeries)
}
rss, isPartial, err := netstorage.ProcessSearchQuery(qt, ec.DenyPartialResponse, sq, ec.Deadline)
rss, isPartial, err := netstorage.ProcessSearchQuery(ec.Context, qt, ec.DenyPartialResponse, sq)
if err != nil {
return nil, err
}
@@ -1947,9 +1948,9 @@ func evalRollupFuncNoCache(qt *querytracer.Tracer, ec *EvalConfig, funcName stri
// Evaluate rollup
keepMetricNames := getKeepMetricNames(expr)
if iafc != nil {
return evalRollupWithIncrementalAggregate(qt, funcName, keepMetricNames, iafc, rss, rcs, preFunc, sharedTimestamps)
return evalRollupWithIncrementalAggregate(ec.Context, qt, funcName, keepMetricNames, iafc, rss, rcs, preFunc, sharedTimestamps)
}
return evalRollupNoIncrementalAggregate(qt, funcName, keepMetricNames, rss, rcs, preFunc, sharedTimestamps)
return evalRollupNoIncrementalAggregate(ec.Context, qt, funcName, keepMetricNames, rss, rcs, preFunc, sharedTimestamps)
}
var (
@@ -1972,14 +1973,14 @@ func maxSilenceInterval() int64 {
return d
}
func evalRollupWithIncrementalAggregate(qt *querytracer.Tracer, funcName string, keepMetricNames bool,
func evalRollupWithIncrementalAggregate(ctx searchutil.Context, qt *querytracer.Tracer, funcName string, keepMetricNames bool,
iafc *incrementalAggrFuncContext, rss *netstorage.Results, rcs []*rollupConfig,
preFunc func(values []float64, timestamps []int64), sharedTimestamps []int64,
) ([]*timeseries, error) {
qt = qt.NewChild("rollup %s() with incremental aggregation %s() over %d series; rollupConfigs=%s", funcName, iafc.ae.Name, rss.Len(), rcs)
defer qt.Done()
var samplesScannedTotal atomic.Uint64
err := rss.RunParallel(qt, func(rs *netstorage.Result, workerID uint) error {
err := rss.RunParallel(ctx, qt, func(rs *netstorage.Result, workerID uint) error {
rs.Values, rs.Timestamps = dropStaleNaNs(funcName, rs.Values, rs.Timestamps)
preFunc(rs.Values, rs.Timestamps)
ts := getTimeseries()
@@ -2013,7 +2014,7 @@ func evalRollupWithIncrementalAggregate(qt *querytracer.Tracer, funcName string,
return tss, nil
}
func evalRollupNoIncrementalAggregate(qt *querytracer.Tracer, funcName string, keepMetricNames bool, rss *netstorage.Results, rcs []*rollupConfig,
func evalRollupNoIncrementalAggregate(ctx searchutil.Context, qt *querytracer.Tracer, funcName string, keepMetricNames bool, rss *netstorage.Results, rcs []*rollupConfig,
preFunc func(values []float64, timestamps []int64), sharedTimestamps []int64,
) ([]*timeseries, error) {
qt = qt.NewChild("rollup %s() over %d series; rollupConfigs=%s", funcName, rss.Len(), rcs)
@@ -2023,7 +2024,7 @@ func evalRollupNoIncrementalAggregate(qt *querytracer.Tracer, funcName string, k
tsw := getTimeseriesByWorkerID()
seriesByWorkerID := tsw.byWorkerID
seriesLen := rss.Len()
err := rss.RunParallel(qt, func(rs *netstorage.Result, workerID uint) error {
err := rss.RunParallel(ctx, qt, func(rs *netstorage.Result, workerID uint) error {
rs.Values, rs.Timestamps = dropStaleNaNs(funcName, rs.Values, rs.Timestamps)
preFunc(rs.Values, rs.Timestamps)
for _, rc := range rcs {

View File

@@ -1,6 +1,7 @@
package promql
import (
"context"
"math"
"testing"
"time"
@@ -65,6 +66,7 @@ func TestExecSuccess(t *testing.T) {
f := func(q string, resultExpected []netstorage.Result) {
t.Helper()
ec := &EvalConfig{
Context: searchutil.NewContext(context.Background(), searchutil.NewDeadline(time.Now(), time.Minute, "")),
AuthTokens: []*auth.Token{{
AccountID: accountID,
ProjectID: projectID,
@@ -75,7 +77,6 @@ func TestExecSuccess(t *testing.T) {
Step: step,
MaxPointsPerSeries: 1e4,
MaxSeries: 1000,
Deadline: searchutil.NewDeadline(time.Now(), time.Minute, ""),
RoundDigits: 100,
}
for range 5 {
@@ -10467,6 +10468,7 @@ func TestExecError(t *testing.T) {
f := func(q string) {
t.Helper()
ec := &EvalConfig{
Context: searchutil.NewContext(context.Background(), searchutil.NewDeadline(time.Now(), time.Minute, "")),
AuthTokens: []*auth.Token{{
AccountID: 123,
ProjectID: 567,
@@ -10476,7 +10478,6 @@ func TestExecError(t *testing.T) {
Step: 100,
MaxPointsPerSeries: 1e4,
MaxSeries: 1000,
Deadline: searchutil.NewDeadline(time.Now(), time.Minute, ""),
RoundDigits: 100,
}
for range 4 {

View File

@@ -1,6 +1,7 @@
package searchutil
import (
"context"
"flag"
"fmt"
"net/http"
@@ -38,34 +39,39 @@ func GetMaxQueryDuration(r *http.Request) time.Duration {
return d
}
// GetDeadlineForQuery returns deadline for the given query r.
func GetDeadlineForQuery(r *http.Request, startTime time.Time) Deadline {
// GetDeadlineForQuery returns context for the given query r.
func GetContextForQuery(r *http.Request, startTime time.Time) Context {
dMax := maxQueryDuration.Milliseconds()
return getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxQueryDuration")
deadline := getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxQueryDuration")
return NewContext(r.Context(), deadline)
}
// GetDeadlineForStatusRequest returns deadline for the given request to /api/v1/status/*.
func GetDeadlineForStatusRequest(r *http.Request, startTime time.Time) Deadline {
// GetContextForStatusRequest returns context for the given request to /api/v1/status/*.
func GetContextForStatusRequest(r *http.Request, startTime time.Time) Context {
dMax := maxStatusRequestDuration.Milliseconds()
return getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxStatusRequestDuration")
deadline := getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxStatusRequestDuration")
return NewContext(r.Context(), deadline)
}
// GetDeadlineForExport returns deadline for the given request to /api/v1/export.
func GetDeadlineForExport(r *http.Request, startTime time.Time) Deadline {
// GetContextForExport returns context for the given request to /api/v1/export.
func GetContextForExport(r *http.Request, startTime time.Time) Context {
dMax := maxExportDuration.Milliseconds()
return getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxExportDuration")
deadline := getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxExportDuration")
return NewContext(r.Context(), deadline)
}
// GetDeadlineForLabelsAPI returns deadline for the given request to /api/v1/labels, /api/v1/label/.../values or /api/v1/series
func GetDeadlineForLabelsAPI(r *http.Request, startTime time.Time) Deadline {
// GetContextForLabelsAPI returns context for the given request to /api/v1/labels, /api/v1/label/.../values or /api/v1/series
func GetContextForLabelsAPI(r *http.Request, startTime time.Time) Context {
dMax := maxLabelsAPIDuration.Milliseconds()
return getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxLabelsAPIDuration")
deadline := getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxLabelsAPIDuration")
return NewContext(r.Context(), deadline)
}
// GetDeadlineForDelete returns deadline for the given request to /api/v1/admin/tsdb/delete_series.
func GetDeadlineForDelete(r *http.Request, startTime time.Time) Deadline {
// GetDeadlineForDelete returns context for the given request to /api/v1/admin/tsdb/delete_series.
func GetContextForDelete(r *http.Request, startTime time.Time) Context {
dMax := maxDeleteDuration.Milliseconds()
return getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxDeleteDuration")
deadline := getDeadlineWithMaxDuration(r, startTime, dMax, "-search.maxDeleteDuration")
return NewContext(r.Context(), deadline)
}
func getDeadlineWithMaxDuration(r *http.Request, startTime time.Time, dMax int64, flagHint string) Deadline {
@@ -80,6 +86,68 @@ func getDeadlineWithMaxDuration(r *http.Request, startTime time.Time, dMax int64
return NewDeadline(startTime, timeout, flagHint)
}
// Context defines search context with deadline
type Context struct {
parent context.Context
deadline Deadline
}
// NewContext return new context for given parent context and deadline
func NewContext(ctx context.Context, deadline Deadline) Context {
if ctx == nil {
ctx = context.Background()
}
return Context{
parent: ctx,
deadline: deadline,
}
}
// NewContextWithDeadlineTimestamp return new context for given parent context and timestamp of deadline
func NewContextWithDeadlineTimestamp(ctx context.Context, timestamp uint64) Context {
deadline := DeadlineFromTimestamp(timestamp)
return Context{
parent: ctx,
deadline: deadline,
}
}
// Deadline returns context deadline
func (ctx *Context) Deadline() Deadline {
return ctx.deadline
}
// IsDone returns true if context is cancelled or deadline exceeded
func (ctx *Context) IsDone() bool {
if ctx.deadline.Exceeded() {
return true
}
if ctx.canceled() {
return true
}
return false
}
func (ctx *Context) canceled() bool {
select {
case <-ctx.parent.Done():
return true
default:
return false
}
}
// Err return context error if there is any
func (ctx *Context) Err() error {
if ctx.deadline.Exceeded() {
return fmt.Errorf("context deadline timeout: %s: %w", ctx.deadline.String(), context.DeadlineExceeded)
}
if ctx.canceled() {
return ctx.parent.Err()
}
return nil
}
// Deadline contains deadline with the corresponding timeout for pretty error messages.
type Deadline struct {
deadline uint64

View File

@@ -276,27 +276,31 @@ func tagFiltersToString(tfs []storage.TagFilter) string {
return string(b)
}
func TestGetDeadline(t *testing.T) {
f := func(got, exp Deadline) {
if got.Deadline() != exp.Deadline() {
t.Fatalf("expected to have %v; got %v instead", exp, got)
func TestGetContextDeadline(t *testing.T) {
f := func(got Context, exp Deadline) {
t.Helper()
// got is a function parameter and therefore addressable,
// so the pointer-receiver Deadline() method resolves here.
gotDeadline := got.Deadline()
if gotDeadline.Deadline() != exp.Deadline() {
t.Fatalf("expected deadline %d; got %d instead", exp.Deadline(), gotDeadline.Deadline())
}
}
start := time.Now()
expDeadline := func(deadline time.Duration) Deadline {
return NewDeadline(start, deadline, "")
expDeadline := func(d time.Duration) Deadline {
return NewDeadline(start, d, "")
}
r, _ := http.NewRequest("GET", "", nil)
f(GetDeadlineForExport(r, start), expDeadline(*maxExportDuration))
f(GetDeadlineForLabelsAPI(r, start), expDeadline(*maxLabelsAPIDuration))
f(GetDeadlineForStatusRequest(r, start), expDeadline(*maxStatusRequestDuration))
f(GetDeadlineForQuery(r, start), expDeadline(*maxQueryDuration))
f(GetContextForExport(r, start), expDeadline(*maxExportDuration))
f(GetContextForLabelsAPI(r, start), expDeadline(*maxLabelsAPIDuration))
f(GetContextForStatusRequest(r, start), expDeadline(*maxStatusRequestDuration))
f(GetContextForQuery(r, start), expDeadline(*maxQueryDuration))
r, _ = http.NewRequest("GET", "http://foo?timeout=1s", nil)
f(GetDeadlineForExport(r, start), expDeadline(time.Second))
f(GetDeadlineForLabelsAPI(r, start), expDeadline(time.Second))
f(GetDeadlineForStatusRequest(r, start), expDeadline(time.Second))
f(GetDeadlineForQuery(r, start), expDeadline(time.Second))
f(GetContextForExport(r, start), expDeadline(time.Second))
f(GetContextForLabelsAPI(r, start), expDeadline(time.Second))
f(GetContextForStatusRequest(r, start), expDeadline(time.Second))
f(GetContextForQuery(r, start), expDeadline(time.Second))
}

View File

@@ -44,7 +44,7 @@ func MetricNamesStatsHandler(startTime time.Time, at *auth.Token, qt *querytrace
return fmt.Errorf("match_pattern=%q must be valid regex: %w", matchPattern, err)
}
}
deadline := searchutil.GetDeadlineForStatusRequest(r, startTime)
ctx := searchutil.GetContextForStatusRequest(r, startTime)
var tt *storage.TenantToken
if at != nil {
tt = &storage.TenantToken{
@@ -52,7 +52,7 @@ func MetricNamesStatsHandler(startTime time.Time, at *auth.Token, qt *querytrace
ProjectID: at.ProjectID,
}
}
stats, err := netstorage.GetMetricNamesStats(qt, tt, limit, le, matchPattern, deadline)
stats, err := netstorage.GetMetricNamesStats(ctx, qt, tt, limit, le, matchPattern)
if err != nil {
return err
}
@@ -62,8 +62,8 @@ func MetricNamesStatsHandler(startTime time.Time, at *auth.Token, qt *querytrace
// ResetMetricNamesStatsHandler resets metric names usage state
func ResetMetricNamesStatsHandler(startTime time.Time, qt *querytracer.Tracer, r *http.Request) error {
deadline := searchutil.GetDeadlineForStatusRequest(r, startTime)
if err := netstorage.ResetMetricNamesStats(qt, deadline); err != nil {
ctx := searchutil.GetContextForStatusRequest(r, startTime)
if err := netstorage.ResetMetricNamesStats(ctx, qt); err != nil {
return err
}
return nil

View File

@@ -26,6 +26,8 @@ See also [LTS releases](https://docs.victoriametrics.com/victoriametrics/lts-rel
## tip
* FEATURE: [vmsingle](https://docs.victoriametrics.com/victoriametrics/single-server-victoriametrics/) and `vmselect` in [VictoriaMetrics cluster](https://docs.victoriametrics.com/victoriametrics/cluster-victoriametrics/): cancel query requests if client closes connection. See [#11355](https://github.com/VictoriaMetrics/VictoriaMetrics/issues/11355).
## [v1.150.0](https://github.com/VictoriaMetrics/VictoriaMetrics/releases/tag/v1.150.0)
Release candidate