Query caching: add request deduplication middleware (#110892)
* secrets: update test to accept []byte(nil) and []byte{} (#110630)
Co-authored-by: Matheus Macabu <macabu.matheus@gmail.com>
* Query caching: add request deduplication middleware
* log error if unable to build cache key
* remove TODO
* always use req.PluginContext.DataSourceInstanceSettings.UID
* make update-workspace
---------
Co-authored-by: Matheus Macabu <macabu.matheus@gmail.com>
This commit is contained in:
@@ -1,7 +1,13 @@
|
||||
package caching
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"encoding/json"
|
||||
"io"
|
||||
"strings"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
)
|
||||
@@ -62,3 +68,68 @@ func (s *OSSCachingService) HandleResourceRequest(ctx context.Context, req *back
|
||||
}
|
||||
|
||||
var _ CachingService = &OSSCachingService{}
|
||||
|
||||
// GetKey creates a prefixed cache key and uses the internal `encoder` to encode the query into a string
|
||||
func GetKey(prefix string, query interface{}) (string, error) {
|
||||
keybuf := bytes.NewBuffer(nil)
|
||||
|
||||
encoder := &JSONEncoder{}
|
||||
|
||||
if err := encoder.Encode(keybuf, query); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
key, err := SHA256KeyFunc(keybuf)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
return strings.Join([]string{prefix, key}, ":"), nil
|
||||
}
|
||||
|
||||
// SHA256KeyFunc copies the data from `r` into a sha256.Hash, and returns the encoded Sum.
|
||||
func SHA256KeyFunc(r io.Reader) (string, error) {
|
||||
hash := sha256.New()
|
||||
|
||||
// Read all data from the provided reader
|
||||
if _, err := io.Copy(hash, r); err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// Encode the written values to SHA256
|
||||
return hex.EncodeToString(hash.Sum(nil)), nil
|
||||
}
|
||||
|
||||
// JSONEncoder encodes and decodes struct data to/from JSON
|
||||
type JSONEncoder struct{}
|
||||
|
||||
// NewJSONEncoder creates a pointer to a new JSONEncoder, which implements the `Encoder` interface
|
||||
func NewJSONEncoder() *JSONEncoder {
|
||||
return &JSONEncoder{}
|
||||
}
|
||||
|
||||
func (e *JSONEncoder) EncodeBytes(w io.Writer, b []byte) error {
|
||||
_, err := w.Write(b)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (e *JSONEncoder) DecodeBytes(r io.Reader) ([]byte, error) {
|
||||
encBytes, err := io.ReadAll(r)
|
||||
if err != nil {
|
||||
return []byte{}, err
|
||||
}
|
||||
return encBytes, err
|
||||
}
|
||||
|
||||
// Encode encodes the `v` interface into `w` using a json.Encoder
|
||||
func (e *JSONEncoder) Encode(w io.Writer, v interface{}) error {
|
||||
return json.NewEncoder(w).Encode(v)
|
||||
}
|
||||
|
||||
// Decode encodes the io.Reader `r` into the interface `v` using a json.Decoder
|
||||
func (e *JSONEncoder) Decode(r io.Reader, v interface{}) error {
|
||||
return json.NewDecoder(r).Decode(v)
|
||||
}
|
||||
|
||||
@@ -339,6 +339,13 @@ var (
|
||||
Expression: "true", // enabled by default
|
||||
Owner: awsDatasourcesSquad,
|
||||
},
|
||||
{
|
||||
Name: "queryCacheRequestDeduplication",
|
||||
Description: "Enable request deduplication when query caching is enabled. Requests issuing the same query will be deduplicated, only the first request to arrive will be executed and the response will be shared with requests arriving while there is a request in-flight",
|
||||
Stage: FeatureStageExperimental,
|
||||
Expression: "false", // enabled by default
|
||||
Owner: grafanaOperatorExperienceSquad,
|
||||
},
|
||||
{
|
||||
Name: "permissionsFilterRemoveSubquery",
|
||||
Description: "Alternative permission filter implementation that does not use subqueries for fetching the dashboard folder",
|
||||
|
||||
@@ -43,6 +43,7 @@ provisioning,experimental,@grafana/grafana-app-platform-squad,false,true,false
|
||||
grafanaAPIServerEnsureKubectlAccess,experimental,@grafana/grafana-app-platform-squad,true,true,false
|
||||
featureToggleAdminPage,experimental,@grafana/grafana-operator-experience-squad,false,true,false
|
||||
awsAsyncQueryCaching,GA,@grafana/aws-datasources,false,false,false
|
||||
queryCacheRequestDeduplication,experimental,@grafana/grafana-operator-experience-squad,false,false,false
|
||||
permissionsFilterRemoveSubquery,experimental,@grafana/search-and-storage,false,false,false
|
||||
configurableSchedulerTick,experimental,@grafana/alerting-squad,false,true,false
|
||||
dashgpt,GA,@grafana/dashboards-squad,false,false,true
|
||||
|
||||
|
@@ -183,6 +183,10 @@ const (
|
||||
// Enable caching for async queries for Redshift and Athena. Requires that the datasource has caching and async query support enabled
|
||||
FlagAwsAsyncQueryCaching = "awsAsyncQueryCaching"
|
||||
|
||||
// FlagQueryCacheRequestDeduplication
|
||||
// Enable request deduplication when query caching is enabled. Requests issuing the same query will be deduplicated, only the first request to arrive will be executed and the response will be shared with requests arriving while there is a request in-flight
|
||||
FlagQueryCacheRequestDeduplication = "queryCacheRequestDeduplication"
|
||||
|
||||
// FlagPermissionsFilterRemoveSubquery
|
||||
// Alternative permission filter implementation that does not use subqueries for fetching the dashboard folder
|
||||
FlagPermissionsFilterRemoveSubquery = "permissionsFilterRemoveSubquery"
|
||||
|
||||
@@ -2887,6 +2887,19 @@
|
||||
"expression": "true"
|
||||
}
|
||||
},
|
||||
{
|
||||
"metadata": {
|
||||
"name": "queryCacheRequestDeduplication",
|
||||
"resourceVersion": "1757521912495",
|
||||
"creationTimestamp": "2025-09-10T16:31:52Z"
|
||||
},
|
||||
"spec": {
|
||||
"description": "Enable request deduplication when query caching is enabled. Requests issuing the same query will be deduplicated, only the first request to arrive will be executed and the response will be shared with requests arriving while there is a request in-flight",
|
||||
"stage": "experimental",
|
||||
"codeowner": "@grafana/grafana-operator-experience-squad",
|
||||
"expression": "false"
|
||||
}
|
||||
},
|
||||
{
|
||||
"metadata": {
|
||||
"name": "queryLibrary",
|
||||
|
||||
@@ -2,12 +2,14 @@ package clientmiddleware
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strconv"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-aws-sdk/pkg/awsds"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/prometheus/client_golang/prometheus"
|
||||
"golang.org/x/sync/singleflight"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/grafana/grafana/pkg/services/caching"
|
||||
@@ -34,14 +36,20 @@ func NewCachingMiddlewareWithFeatureManager(cachingService caching.CachingServic
|
||||
if err := prometheus.Register(ResourceCachingRequestHistogram); err != nil {
|
||||
log.Error("Error registering prometheus collector 'ResourceRequestHistogram'", "error", err)
|
||||
}
|
||||
return backend.HandlerMiddlewareFunc(func(next backend.Handler) backend.Handler {
|
||||
return &CachingMiddleware{
|
||||
cachingMiddlewareHandler := func(next backend.Handler) backend.Handler {
|
||||
cachingMiddleware := &CachingMiddleware{
|
||||
BaseHandler: backend.NewBaseHandler(next),
|
||||
caching: cachingService,
|
||||
log: log,
|
||||
features: features,
|
||||
}
|
||||
})
|
||||
if features != nil && features.IsEnabled(context.Background(), featuremgmt.FlagQueryCacheRequestDeduplication) {
|
||||
return newRequestDeduplicationMiddleware(log, cachingMiddleware)
|
||||
}
|
||||
return cachingMiddleware
|
||||
}
|
||||
|
||||
return backend.HandlerMiddlewareFunc(cachingMiddlewareHandler)
|
||||
}
|
||||
|
||||
type CachingMiddleware struct {
|
||||
@@ -164,3 +172,52 @@ func (m *CachingMiddleware) CallResource(ctx context.Context, req *backend.CallR
|
||||
|
||||
return m.BaseHandler.CallResource(ctx, req, cacheSender)
|
||||
}
|
||||
|
||||
// Given N requests happening at the same time and issuing the same query, only one request will execute
|
||||
// and the other ones will wait for the response received by the request being executed.
|
||||
type requestDeduplicationMiddleware struct {
|
||||
backend.BaseHandler
|
||||
log *log.ConcreteLogger
|
||||
singleflight *singleflight.Group
|
||||
}
|
||||
|
||||
func newRequestDeduplicationMiddleware(log *log.ConcreteLogger, next backend.Handler) *requestDeduplicationMiddleware {
|
||||
return &requestDeduplicationMiddleware{log: log, BaseHandler: backend.NewBaseHandler(next), singleflight: &singleflight.Group{}}
|
||||
}
|
||||
|
||||
func (m *requestDeduplicationMiddleware) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
if req.PluginContext.DataSourceInstanceSettings.UID == "" {
|
||||
return m.BaseHandler.QueryData(ctx, req)
|
||||
}
|
||||
key, err := caching.GetKey(req.PluginContext.DataSourceInstanceSettings.UID, req)
|
||||
if err != nil {
|
||||
m.log.Error("error building cache key for request deduplication, skipping request deduplication", "error", err)
|
||||
return m.BaseHandler.QueryData(ctx, req)
|
||||
}
|
||||
v, err, _ := m.singleflight.Do(key, func() (interface{}, error) {
|
||||
return m.BaseHandler.QueryData(ctx, req)
|
||||
})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("singleflight: calling next.QueryData: %w", err)
|
||||
}
|
||||
return v.(*backend.QueryDataResponse), nil
|
||||
}
|
||||
|
||||
func (m *requestDeduplicationMiddleware) CallResource(ctx context.Context, req *backend.CallResourceRequest, sender backend.CallResourceResponseSender) error {
|
||||
if req.PluginContext.DataSourceInstanceSettings.UID == "" {
|
||||
return m.BaseHandler.CallResource(ctx, req, sender)
|
||||
}
|
||||
|
||||
key, err := caching.GetKey(req.PluginContext.DataSourceInstanceSettings.UID, req)
|
||||
if err != nil {
|
||||
m.log.Error("error building cache key for request deduplication, skipping request deduplication", "error", err)
|
||||
return m.BaseHandler.CallResource(ctx, req, sender)
|
||||
}
|
||||
_, err, _ = m.singleflight.Do(key, func() (interface{}, error) {
|
||||
return nil, m.BaseHandler.CallResource(ctx, req, sender)
|
||||
})
|
||||
if err != nil {
|
||||
return fmt.Errorf("singleflight: calling next.CallResource: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -4,7 +4,10 @@ import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend"
|
||||
"github.com/grafana/grafana-plugin-sdk-go/backend/handlertest"
|
||||
@@ -322,3 +325,51 @@ func TestCachingMiddleware(t *testing.T) {
|
||||
})
|
||||
})
|
||||
}
|
||||
|
||||
func TestRequestDeduplicationMiddleware(t *testing.T) {
|
||||
t.Run("deduplicates requests issuing the same query", func(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
handler := newMockMiddlewareHandler()
|
||||
middleware := newRequestDeduplicationMiddleware(nil, handler)
|
||||
|
||||
req := backend.QueryDataRequest{
|
||||
PluginContext: backend.PluginContext{
|
||||
DataSourceInstanceSettings: &backend.DataSourceInstanceSettings{
|
||||
UID: "uid",
|
||||
},
|
||||
},
|
||||
}
|
||||
|
||||
wg := &sync.WaitGroup{}
|
||||
wg.Add(2)
|
||||
|
||||
for range 2 {
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
resp, err := middleware.QueryData(t.Context(), &req)
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, &backend.QueryDataResponse{}, resp)
|
||||
}()
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
require.EqualValues(t, 1, handler.QueryDataCalls)
|
||||
})
|
||||
}
|
||||
|
||||
type mockMiddlewareHandler struct {
|
||||
backend.BaseHandler
|
||||
QueryDataCalls int32
|
||||
}
|
||||
|
||||
func newMockMiddlewareHandler() *mockMiddlewareHandler {
|
||||
return &mockMiddlewareHandler{}
|
||||
}
|
||||
|
||||
func (m *mockMiddlewareHandler) QueryData(ctx context.Context, req *backend.QueryDataRequest) (*backend.QueryDataResponse, error) {
|
||||
atomic.AddInt32(&m.QueryDataCalls, 1)
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
return &backend.QueryDataResponse{}, nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user