diff --git a/internal/api/api_test.go b/internal/api/api_test.go index 85682363..0348c169 100644 --- a/internal/api/api_test.go +++ b/internal/api/api_test.go @@ -212,7 +212,7 @@ func TestRestApi(t *testing.T) { }, }} - metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int) (schema.JobData, error) { + metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int, resampleAlgo string) (schema.JobData, error) { return testData, nil } @@ -513,7 +513,7 @@ func TestStopJobWithReusedJobId(t *testing.T) { }, }} - metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int) (schema.JobData, error) { + metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int, resampleAlgo string) (schema.JobData, error) { return testData, nil } diff --git a/internal/api/job.go b/internal/api/job.go index 30ce60ec..3d6348b4 100644 --- a/internal/api/job.go +++ b/internal/api/job.go @@ -309,7 +309,7 @@ func (api *RestAPI) getCompleteJobByID(rw http.ResponseWriter, r *http.Request) } if r.URL.Query().Get("all-metrics") == "true" { - data, err = metricdispatch.LoadData(job, nil, scopes, r.Context(), resolution, "") + data, err = metricdispatch.LoadData(job, nil, scopes, r.Context(), resolution, config.ResampleAlgo()) if err != nil { cclog.Warnf("REST: error while loading all-metrics job data for JobID %d on %s", job.JobID, job.Cluster) return @@ -405,7 +405,7 @@ func (api *RestAPI) getJobByID(rw http.ResponseWriter, r *http.Request) { resolution = max(resolution, mc.Timestep) } - data, err := metricdispatch.LoadData(job, metrics, scopes, r.Context(), resolution, "") + data, err := metricdispatch.LoadData(job, metrics, scopes, r.Context(), resolution, config.ResampleAlgo()) if err != nil { cclog.Warnf("REST: error while loading job data for JobID %d on %s", job.JobID, job.Cluster) return diff --git a/internal/api/nats_test.go b/internal/api/nats_test.go index f6d4b5eb..5bffe1f4 100644 --- a/internal/api/nats_test.go +++ b/internal/api/nats_test.go @@ -547,7 +547,7 @@ func TestNatsHandleStopJob(t *testing.T) { }, }} - metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int) (schema.JobData, error) { + metricstore.TestLoadDataCallback = func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int, resampleAlgo string) (schema.JobData, error) { return testData, nil } diff --git a/internal/config/config.go b/internal/config/config.go index 5865c5ec..6230f70f 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -193,6 +193,21 @@ func initResampler() { // DefaultResamplePolicy is used when no resample policy is configured. const DefaultResamplePolicy = "medium" +// DefaultResampleAlgo is used when neither the user nor the config selects a +// resample algorithm. "average" performs RRDTool-style interval averaging, +// which keeps each plotted point a true mean of its interval. +const DefaultResampleAlgo = "average" + +// ResampleAlgo returns the configured default resample algorithm, falling back +// to DefaultResampleAlgo. Note that an empty string would select LTTB in +// cc-lib's resampler, so callers must not pass "" when they mean "the default". +func ResampleAlgo() string { + if Keys.EnableResampling != nil && Keys.EnableResampling.DefaultAlgo != "" { + return Keys.EnableResampling.DefaultAlgo + } + return DefaultResampleAlgo +} + // TargetPointsForPolicy returns the target number of data points for a resample // policy. This is the single source of truth: it feeds both the requested // resolution (via metricdispatch.ComputeResolution) and the resampler's diff --git a/internal/metricdispatch/dataLoader.go b/internal/metricdispatch/dataLoader.go index 16b6be4e..7736202e 100644 --- a/internal/metricdispatch/dataLoader.go +++ b/internal/metricdispatch/dataLoader.go @@ -23,13 +23,14 @@ // - Running jobs: 2 minutes (data changes frequently) // - Completed jobs: 5 hours (data is static) // -// The cache key is based on job ID, state, requested metrics, scopes, and resolution. +// The cache key is based on job ID, state, requested metrics, scopes, resolution, +// and resample algorithm. // // # Usage // // The primary entry point is LoadData, which automatically handles both running and archived jobs: // -// jobData, err := metricdispatch.LoadData(job, metrics, scopes, ctx, resolution) +// jobData, err := metricdispatch.LoadData(job, metrics, scopes, ctx, resolution, resampleAlgo) // if err != nil { // // Handle error // } @@ -115,7 +116,7 @@ func LoadData(job *schema.Job, } } - jd, err = ms.LoadData(job, metrics, scopes, ctx, resolution) + jd, err = ms.LoadData(job, metrics, scopes, ctx, resolution, resampleAlgo) if err != nil { if len(jd.Metrics) != 0 { cclog.Warnf("partial error loading metrics from store for job %d (user: %s, project: %s, cluster: %s-%s): %s", diff --git a/internal/metricdispatch/metricdata.go b/internal/metricdispatch/metricdata.go index 41e702bd..21d839aa 100755 --- a/internal/metricdispatch/metricdata.go +++ b/internal/metricdispatch/metricdata.go @@ -24,7 +24,8 @@ type MetricDataRepository interface { metrics []string, scopes []schema.MetricScope, ctx context.Context, - resolution int) (schema.JobData, error) + resolution int, + resampleAlgo string) (schema.JobData, error) // Return a map of metrics to a map of nodes to the metric statistics of the job. node scope only. LoadStats(job *schema.Job, diff --git a/internal/metricstoreclient/cc-metric-store.go b/internal/metricstoreclient/cc-metric-store.go index 11132a08..db227bb0 100644 --- a/internal/metricstoreclient/cc-metric-store.go +++ b/internal/metricstoreclient/cc-metric-store.go @@ -81,13 +81,14 @@ type CCMetricStore struct { // APIQueryRequest represents a request to the cc-metric-store query API. // It supports both explicit queries and "for-all-nodes" bulk queries. type APIQueryRequest struct { - Cluster string `json:"cluster"` // Target cluster name - Queries []APIQuery `json:"queries"` // Explicit list of metric queries - ForAllNodes []string `json:"for-all-nodes"` // Metrics to query for all nodes - From int64 `json:"from"` // Start time (Unix timestamp) - To int64 `json:"to"` // End time (Unix timestamp) - WithStats bool `json:"with-stats"` // Include min/avg/max statistics - WithData bool `json:"with-data"` // Include time series data points + Cluster string `json:"cluster"` // Target cluster name + Queries []APIQuery `json:"queries"` // Explicit list of metric queries + ForAllNodes []string `json:"for-all-nodes"` // Metrics to query for all nodes + From int64 `json:"from"` // Start time (Unix timestamp) + To int64 `json:"to"` // End time (Unix timestamp) + WithStats bool `json:"with-stats"` // Include min/avg/max statistics + WithData bool `json:"with-data"` // Include time series data points + ResampleAlgo string `json:"resample-algo,omitempty"` // Downsampling algorithm ("lttb", "average", "simple"); empty = server default } // APIQuery specifies a single metric query with optional scope filtering. @@ -231,6 +232,7 @@ func (ccms *CCMetricStore) LoadData( scopes []schema.MetricScope, ctx context.Context, resolution int, + resampleAlgo string, ) (schema.JobData, error) { queries, assignedScope, err := ccms.buildQueries(job, metrics, scopes, resolution) if err != nil { @@ -245,12 +247,13 @@ func (ccms *CCMetricStore) LoadData( } req := APIQueryRequest{ - Cluster: job.Cluster, - From: job.StartTime, - To: job.StartTime + int64(job.Duration), - Queries: queries, - WithStats: true, - WithData: true, + Cluster: job.Cluster, + From: job.StartTime, + To: job.StartTime + int64(job.Duration), + Queries: queries, + WithStats: true, + WithData: true, + ResampleAlgo: resampleAlgo, } resBody, err := ccms.doRequest(ctx, &req) @@ -632,12 +635,13 @@ func (ccms *CCMetricStore) LoadNodeListData( } req := APIQueryRequest{ - Cluster: cluster, - Queries: queries, - From: from.Unix(), - To: to.Unix(), - WithStats: true, - WithData: true, + Cluster: cluster, + Queries: queries, + From: from.Unix(), + To: to.Unix(), + WithStats: true, + WithData: true, + ResampleAlgo: resampleAlgo, } resBody, err := ccms.doRequest(ctx, &req) diff --git a/pkg/metricstore/query.go b/pkg/metricstore/query.go index 81edd7dc..650025a2 100644 --- a/pkg/metricstore/query.go +++ b/pkg/metricstore/query.go @@ -50,7 +50,7 @@ func (ccms *InternalMetricStore) HealthCheck(cluster string, // TestLoadDataCallback allows tests to override LoadData behavior for testing purposes. // When set to a non-nil function, LoadData will call this function instead of the default implementation. -var TestLoadDataCallback func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int) (schema.JobData, error) +var TestLoadDataCallback func(job *schema.Job, metrics []string, scopes []schema.MetricScope, ctx context.Context, resolution int, resampleAlgo string) (schema.JobData, error) // LoadData loads metric data for a specific job with automatic scope transformation. // @@ -81,9 +81,10 @@ func (ccms *InternalMetricStore) LoadData( scopes []schema.MetricScope, ctx context.Context, resolution int, + resampleAlgo string, ) (schema.JobData, error) { if TestLoadDataCallback != nil { - return TestLoadDataCallback(job, metrics, scopes, ctx, resolution) + return TestLoadDataCallback(job, metrics, scopes, ctx, resolution, resampleAlgo) } queries, assignedScope, err := buildQueries(job, metrics, scopes, int64(resolution)) @@ -99,12 +100,13 @@ func (ccms *InternalMetricStore) LoadData( } req := APIQueryRequest{ - Cluster: job.Cluster, - From: job.StartTime, - To: job.StartTime + int64(job.Duration), - Queries: queries, - WithStats: true, - WithData: true, + Cluster: job.Cluster, + From: job.StartTime, + To: job.StartTime + int64(job.Duration), + Queries: queries, + WithStats: true, + WithData: true, + ResampleAlgo: resampleAlgo, } resBody, err := FetchData(req)