feat(dashboard): Org-aware cache for schema migration (#115025)
* fix: use dsIndexProvider cache on migrations * chore: use same comment as before * feat: org-aware TTL cache for schemaversion migration and warmup for single tenant * chore: use LRU cache * chore: change DefaultCacheTTL to 1 minute * chore: address copilot reviews * chore: use expirable cache * chore: remove unused import
This commit is contained in:
@@ -0,0 +1,104 @@
|
||||
package schemaversion
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/authlib/types"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/hashicorp/golang-lru/v2/expirable"
|
||||
k8srequest "k8s.io/apiserver/pkg/endpoints/request"
|
||||
|
||||
"github.com/grafana/grafana/pkg/services/apiserver/endpoints/request"
|
||||
)
|
||||
|
||||
const defaultCacheSize = 1000
|
||||
|
||||
// CacheProvider is a generic cache interface for schema version providers.
|
||||
type CacheProvider[T any] interface {
|
||||
// Get returns the cached value if it's still valid, otherwise calls fetch and caches the result.
|
||||
Get(ctx context.Context) T
|
||||
}
|
||||
|
||||
// PreloadableCache is an interface for providers that support preloading the cache.
|
||||
type PreloadableCache interface {
|
||||
// Preload loads data into the cache for the given namespaces.
|
||||
Preload(ctx context.Context, nsInfos []types.NamespaceInfo)
|
||||
}
|
||||
|
||||
// cachedProvider is a thread-safe TTL cache that wraps any fetch function.
|
||||
type cachedProvider[T any] struct {
|
||||
fetch func(context.Context) T
|
||||
cache *expirable.LRU[string, T] // LRU cache: namespace to cache entry
|
||||
inFlight sync.Map // map[string]*sync.Mutex - per-namespace fetch locks
|
||||
logger log.Logger
|
||||
}
|
||||
|
||||
// newCachedProvider creates a new cachedProvider.
|
||||
// The fetch function should be able to handle context with different namespaces.
|
||||
// A non-positive size turns LRU mechanism off (cache of unlimited size).
|
||||
// A non-positive cacheTTL disables TTL expiration.
|
||||
func newCachedProvider[T any](fetch func(context.Context) T, size int, cacheTTL time.Duration, logger log.Logger) *cachedProvider[T] {
|
||||
cacheProvider := &cachedProvider[T]{
|
||||
fetch: fetch,
|
||||
logger: logger,
|
||||
}
|
||||
cacheProvider.cache = expirable.NewLRU(size, func(key string, value T) {
|
||||
cacheProvider.inFlight.Delete(key)
|
||||
}, cacheTTL)
|
||||
return cacheProvider
|
||||
}
|
||||
|
||||
// Get returns the cached value if it's still valid, otherwise calls fetch and caches the result.
|
||||
func (p *cachedProvider[T]) Get(ctx context.Context) T {
|
||||
// Get namespace info from ctx
|
||||
nsInfo, err := request.NamespaceInfoFrom(ctx, true)
|
||||
if err != nil {
|
||||
// No namespace, fall back to direct fetch call without caching
|
||||
p.logger.Warn("Unable to get namespace info from context, skipping cache", "error", err)
|
||||
return p.fetch(ctx)
|
||||
}
|
||||
|
||||
namespace := nsInfo.Value
|
||||
// Fast path: check if cache is still valid
|
||||
if entry, ok := p.cache.Get(namespace); ok {
|
||||
return entry
|
||||
}
|
||||
|
||||
// Get or create a per-namespace lock for this fetch operation
|
||||
// This ensures only one fetch happens per namespace at a time
|
||||
lockInterface, _ := p.inFlight.LoadOrStore(namespace, &sync.Mutex{})
|
||||
nsMutex := lockInterface.(*sync.Mutex)
|
||||
|
||||
// Lock this specific namespace - other namespaces can still proceed
|
||||
nsMutex.Lock()
|
||||
defer nsMutex.Unlock()
|
||||
|
||||
// Double-check: another goroutine might have already fetched while we waited
|
||||
if entry, ok := p.cache.Get(namespace); ok {
|
||||
return entry
|
||||
}
|
||||
|
||||
// Fetch outside the main lock - only this namespace is blocked
|
||||
p.logger.Debug("cache miss or expired, fetching new value", "namespace", namespace)
|
||||
value := p.fetch(ctx)
|
||||
|
||||
// Update the cache for this namespace
|
||||
p.cache.Add(namespace, value)
|
||||
|
||||
return value
|
||||
}
|
||||
|
||||
// Preload loads data into the cache for the given namespaces.
|
||||
func (p *cachedProvider[T]) Preload(ctx context.Context, nsInfos []types.NamespaceInfo) {
|
||||
// Build the cache using a context with the namespace
|
||||
p.logger.Info("preloading cache", "nsInfos", len(nsInfos))
|
||||
startedAt := time.Now()
|
||||
defer func() {
|
||||
p.logger.Info("finished preloading cache", "nsInfos", len(nsInfos), "elapsed", time.Since(startedAt))
|
||||
}()
|
||||
for _, nsInfo := range nsInfos {
|
||||
p.cache.Add(nsInfo.Value, p.fetch(k8srequest.WithNamespace(ctx, nsInfo.Value)))
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,478 @@
|
||||
package schemaversion
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
authlib "github.com/grafana/authlib/types"
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"k8s.io/apiserver/pkg/endpoints/request"
|
||||
)
|
||||
|
||||
// testProvider tracks how many times get() is called
|
||||
type testProvider struct {
|
||||
testData any
|
||||
callCount atomic.Int64
|
||||
}
|
||||
|
||||
func newTestProvider(testData any) *testProvider {
|
||||
return &testProvider{
|
||||
testData: testData,
|
||||
}
|
||||
}
|
||||
|
||||
func (p *testProvider) get(_ context.Context) any {
|
||||
p.callCount.Add(1)
|
||||
return p.testData
|
||||
}
|
||||
|
||||
func (p *testProvider) getCallCount() int64 {
|
||||
return p.callCount.Load()
|
||||
}
|
||||
|
||||
func TestCachedProvider_CacheHit(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
{UID: "ds2", Type: "loki", Name: "Loki"},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
// Test newCachedProvider directly instead of the wrapper
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Use "default" namespace (org 1) - this is the standard Grafana namespace format
|
||||
ctx := request.WithNamespace(context.Background(), "default")
|
||||
|
||||
// First call should hit the underlying provider
|
||||
idx1 := cached.Get(ctx)
|
||||
require.NotNil(t, idx1)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount(), "first call should invoke underlying provider")
|
||||
|
||||
// Second call should use cache
|
||||
idx2 := cached.Get(ctx)
|
||||
require.NotNil(t, idx2)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount(), "second call should use cache, not invoke underlying provider")
|
||||
|
||||
// Both should return the same data
|
||||
assert.Equal(t, idx1, idx2)
|
||||
}
|
||||
|
||||
func TestCachedProvider_NamespaceIsolation(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Use "default" (org 1) and "org-2" (org 2) - standard Grafana namespace formats
|
||||
ctx1 := request.WithNamespace(context.Background(), "default")
|
||||
ctx2 := request.WithNamespace(context.Background(), "org-2")
|
||||
|
||||
// First call for org 1
|
||||
idx1 := cached.Get(ctx1)
|
||||
require.NotNil(t, idx1)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount(), "first org-1 call should invoke underlying provider")
|
||||
|
||||
// Call for org 2 should also invoke underlying provider (different namespace)
|
||||
idx2 := cached.Get(ctx2)
|
||||
require.NotNil(t, idx2)
|
||||
assert.Equal(t, int64(2), underlying.getCallCount(), "org-2 call should invoke underlying provider (separate cache)")
|
||||
|
||||
// Second call for org 1 should use cache
|
||||
idx3 := cached.Get(ctx1)
|
||||
require.NotNil(t, idx3)
|
||||
assert.Equal(t, int64(2), underlying.getCallCount(), "second org-1 call should use cache")
|
||||
|
||||
// Second call for org 2 should use cache
|
||||
idx4 := cached.Get(ctx2)
|
||||
require.NotNil(t, idx4)
|
||||
assert.Equal(t, int64(2), underlying.getCallCount(), "second org-2 call should use cache")
|
||||
}
|
||||
|
||||
func TestCachedProvider_NoNamespaceFallback(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Context without namespace - should fall back to direct provider call
|
||||
ctx := context.Background()
|
||||
|
||||
idx1 := cached.Get(ctx)
|
||||
require.NotNil(t, idx1)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount())
|
||||
|
||||
// Second call without namespace should also invoke underlying (no caching for unknown namespace)
|
||||
idx2 := cached.Get(ctx)
|
||||
require.NotNil(t, idx2)
|
||||
assert.Equal(t, int64(2), underlying.getCallCount(), "without namespace, each call should invoke underlying provider")
|
||||
}
|
||||
|
||||
func TestCachedProvider_ConcurrentAccess(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Use "default" namespace (org 1)
|
||||
ctx := request.WithNamespace(context.Background(), "default")
|
||||
|
||||
var wg sync.WaitGroup
|
||||
numGoroutines := 100
|
||||
|
||||
// Launch many goroutines that all try to access the cache simultaneously
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
idx := cached.Get(ctx)
|
||||
require.NotNil(t, idx)
|
||||
}()
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
// Due to double-check locking, only 1 goroutine should have actually built the cache
|
||||
// In practice, there might be a few more due to timing, but it should be much less than numGoroutines
|
||||
callCount := underlying.getCallCount()
|
||||
assert.LessOrEqual(t, callCount, int64(5), "with proper locking, very few goroutines should invoke underlying provider; got %d", callCount)
|
||||
}
|
||||
|
||||
func TestCachedProvider_ConcurrentNamespaces(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
var wg sync.WaitGroup
|
||||
numOrgs := 10
|
||||
callsPerOrg := 20
|
||||
|
||||
// Launch goroutines for multiple namespaces
|
||||
// Use valid namespace formats: "default" for org 1, "org-N" for N > 1
|
||||
namespaces := make([]string, numOrgs)
|
||||
namespaces[0] = "default"
|
||||
for i := 1; i < numOrgs; i++ {
|
||||
namespaces[i] = fmt.Sprintf("org-%d", i+1)
|
||||
}
|
||||
|
||||
for _, ns := range namespaces {
|
||||
ctx := request.WithNamespace(context.Background(), ns)
|
||||
for i := 0; i < callsPerOrg; i++ {
|
||||
wg.Add(1)
|
||||
go func(ctx context.Context) {
|
||||
defer wg.Done()
|
||||
idx := cached.Get(ctx)
|
||||
require.NotNil(t, idx)
|
||||
}(ctx)
|
||||
}
|
||||
}
|
||||
|
||||
wg.Wait()
|
||||
|
||||
// Each org should have at most a few calls (ideally 1, but timing can cause a few more)
|
||||
callCount := underlying.getCallCount()
|
||||
// With 10 orgs, we expect around 10 calls (one per org)
|
||||
assert.LessOrEqual(t, callCount, int64(numOrgs), "expected roughly one call per org, got %d calls for %d orgs", callCount, numOrgs)
|
||||
}
|
||||
|
||||
// Test that cache returns correct data for each namespace
|
||||
func TestCachedProvider_CorrectDataPerNamespace(t *testing.T) {
|
||||
// Provider that returns different data based on namespace
|
||||
underlying := &namespaceAwareProvider{
|
||||
datasourcesByNamespace: map[string][]DataSourceInfo{
|
||||
"default": {{UID: "org1-ds", Type: "prometheus", Name: "Org1 DS", Default: true}},
|
||||
"org-2": {{UID: "org2-ds", Type: "loki", Name: "Org2 DS", Default: true}},
|
||||
},
|
||||
}
|
||||
cached := newCachedProvider(underlying.Index, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Use valid namespace formats
|
||||
ctx1 := request.WithNamespace(context.Background(), "default")
|
||||
ctx2 := request.WithNamespace(context.Background(), "org-2")
|
||||
|
||||
idx1 := cached.Get(ctx1)
|
||||
idx2 := cached.Get(ctx2)
|
||||
|
||||
assert.Equal(t, "org1-ds", idx1.GetDefault().UID, "org 1 should get org-1 datasources")
|
||||
assert.Equal(t, "org2-ds", idx2.GetDefault().UID, "org 2 should get org-2 datasources")
|
||||
|
||||
// Subsequent calls should still return correct data
|
||||
idx1Again := cached.Get(ctx1)
|
||||
idx2Again := cached.Get(ctx2)
|
||||
|
||||
assert.Equal(t, "org1-ds", idx1Again.GetDefault().UID, "org 1 should still get org-1 datasources from cache")
|
||||
assert.Equal(t, "org2-ds", idx2Again.GetDefault().UID, "org 2 should still get org-2 datasources from cache")
|
||||
}
|
||||
|
||||
// TestCachedProvider_PreloadMultipleNamespaces verifies preloading multiple namespaces
|
||||
func TestCachedProvider_PreloadMultipleNamespaces(t *testing.T) {
|
||||
// Provider that returns different data based on namespace
|
||||
underlying := &namespaceAwareProvider{
|
||||
datasourcesByNamespace: map[string][]DataSourceInfo{
|
||||
"default": {{UID: "org1-ds", Type: "prometheus", Name: "Org1 DS", Default: true}},
|
||||
"org-2": {{UID: "org2-ds", Type: "loki", Name: "Org2 DS", Default: true}},
|
||||
"org-3": {{UID: "org3-ds", Type: "tempo", Name: "Org3 DS", Default: true}},
|
||||
},
|
||||
}
|
||||
cached := newCachedProvider(underlying.Index, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
// Preload multiple namespaces
|
||||
nsInfos := []authlib.NamespaceInfo{
|
||||
createNamespaceInfo(1, 0, "default"),
|
||||
createNamespaceInfo(2, 0, "org-2"),
|
||||
createNamespaceInfo(3, 0, "org-3"),
|
||||
}
|
||||
cached.Preload(context.Background(), nsInfos)
|
||||
|
||||
// After preload, the underlying provider should have been called once per namespace
|
||||
assert.Equal(t, 3, underlying.callCount, "preload should call underlying provider once per namespace")
|
||||
|
||||
// Access all namespaces - should use preloaded data and get correct data per namespace
|
||||
expectedUIDs := map[string]string{
|
||||
"default": "org1-ds",
|
||||
"org-2": "org2-ds",
|
||||
"org-3": "org3-ds",
|
||||
}
|
||||
|
||||
for _, ns := range []string{"default", "org-2", "org-3"} {
|
||||
ctx := request.WithNamespace(context.Background(), ns)
|
||||
idx := cached.Get(ctx)
|
||||
require.NotNil(t, idx, "index for namespace %s should not be nil", ns)
|
||||
assert.Equal(t, expectedUIDs[ns], idx.GetDefault().UID, "namespace %s should get correct datasource", ns)
|
||||
}
|
||||
|
||||
// The underlying provider should still have been called only 3 times (from preload)
|
||||
assert.Equal(t, 3, underlying.callCount,
|
||||
"access after preload should use cached data for all namespaces")
|
||||
}
|
||||
|
||||
// namespaceAwareProvider returns different datasources based on namespace
|
||||
type namespaceAwareProvider struct {
|
||||
datasourcesByNamespace map[string][]DataSourceInfo
|
||||
callCount int
|
||||
}
|
||||
|
||||
func (p *namespaceAwareProvider) Index(ctx context.Context) *DatasourceIndex {
|
||||
p.callCount++
|
||||
ns := request.NamespaceValue(ctx)
|
||||
if ds, ok := p.datasourcesByNamespace[ns]; ok {
|
||||
return NewDatasourceIndex(ds)
|
||||
}
|
||||
return NewDatasourceIndex(nil)
|
||||
}
|
||||
|
||||
// createNamespaceInfo creates a NamespaceInfo for testing
|
||||
func createNamespaceInfo(orgID, stackID int64, value string) authlib.NamespaceInfo {
|
||||
return authlib.NamespaceInfo{
|
||||
OrgID: orgID,
|
||||
StackID: stackID,
|
||||
Value: value,
|
||||
}
|
||||
}
|
||||
|
||||
// Test DatasourceIndex functionality
|
||||
func TestDatasourceIndex_Lookup(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds-uid-1", Type: "prometheus", Name: "Prometheus DS", Default: true, APIVersion: "v1"},
|
||||
{UID: "ds-uid-2", Type: "loki", Name: "Loki DS", Default: false, APIVersion: "v1"},
|
||||
}
|
||||
idx := NewDatasourceIndex(datasources)
|
||||
|
||||
t.Run("lookup by name", func(t *testing.T) {
|
||||
ds := idx.Lookup("Prometheus DS")
|
||||
require.NotNil(t, ds)
|
||||
assert.Equal(t, "ds-uid-1", ds.UID)
|
||||
})
|
||||
|
||||
t.Run("lookup by UID", func(t *testing.T) {
|
||||
ds := idx.Lookup("ds-uid-2")
|
||||
require.NotNil(t, ds)
|
||||
assert.Equal(t, "Loki DS", ds.Name)
|
||||
})
|
||||
|
||||
t.Run("lookup unknown returns nil", func(t *testing.T) {
|
||||
ds := idx.Lookup("unknown")
|
||||
assert.Nil(t, ds)
|
||||
})
|
||||
|
||||
t.Run("get default", func(t *testing.T) {
|
||||
ds := idx.GetDefault()
|
||||
require.NotNil(t, ds)
|
||||
assert.Equal(t, "ds-uid-1", ds.UID)
|
||||
})
|
||||
|
||||
t.Run("lookup by UID directly", func(t *testing.T) {
|
||||
ds := idx.LookupByUID("ds-uid-1")
|
||||
require.NotNil(t, ds)
|
||||
assert.Equal(t, "Prometheus DS", ds.Name)
|
||||
})
|
||||
|
||||
t.Run("lookup by name directly", func(t *testing.T) {
|
||||
ds := idx.LookupByName("Loki DS")
|
||||
require.NotNil(t, ds)
|
||||
assert.Equal(t, "ds-uid-2", ds.UID)
|
||||
})
|
||||
}
|
||||
|
||||
func TestDatasourceIndex_EmptyIndex(t *testing.T) {
|
||||
idx := NewDatasourceIndex(nil)
|
||||
|
||||
assert.Nil(t, idx.GetDefault())
|
||||
assert.Nil(t, idx.Lookup("anything"))
|
||||
assert.Nil(t, idx.LookupByUID("anything"))
|
||||
assert.Nil(t, idx.LookupByName("anything"))
|
||||
}
|
||||
|
||||
// TestCachedProvider_TTLExpiration verifies that cache expires after TTL
|
||||
func TestCachedProvider_TTLExpiration(t *testing.T) {
|
||||
datasources := []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
}
|
||||
|
||||
underlying := newTestProvider(datasources)
|
||||
// Use a very short TTL for testing
|
||||
shortTTL := 50 * time.Millisecond
|
||||
cached := newCachedProvider(underlying.get, defaultCacheSize, shortTTL, log.New("test"))
|
||||
|
||||
ctx := request.WithNamespace(context.Background(), "default")
|
||||
|
||||
// First call - should call underlying provider
|
||||
idx1 := cached.Get(ctx)
|
||||
require.NotNil(t, idx1)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount(), "first call should invoke underlying provider")
|
||||
|
||||
// Second call immediately - should use cache
|
||||
idx2 := cached.Get(ctx)
|
||||
require.NotNil(t, idx2)
|
||||
assert.Equal(t, int64(1), underlying.getCallCount(), "second call should use cache")
|
||||
|
||||
// Wait for TTL to expire
|
||||
time.Sleep(shortTTL + 20*time.Millisecond)
|
||||
|
||||
// Third call after TTL - should call underlying provider again
|
||||
idx3 := cached.Get(ctx)
|
||||
require.NotNil(t, idx3)
|
||||
assert.Equal(t, int64(2), underlying.getCallCount(),
|
||||
"after TTL expiration, underlying provider should be called again")
|
||||
}
|
||||
|
||||
// TestCachedProvider_ParallelNamespacesFetch verifies that different namespaces can fetch in parallel
|
||||
func TestCachedProvider_ParallelNamespacesFetch(t *testing.T) {
|
||||
// Create a blocking provider that tracks concurrent executions
|
||||
provider := &blockingProvider{
|
||||
blockDuration: 100 * time.Millisecond,
|
||||
datasources: []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
},
|
||||
}
|
||||
cached := newCachedProvider(provider.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
numNamespaces := 5
|
||||
var wg sync.WaitGroup
|
||||
|
||||
// Launch fetches for different namespaces simultaneously
|
||||
startTime := time.Now()
|
||||
for i := 0; i < numNamespaces; i++ {
|
||||
wg.Add(1)
|
||||
namespace := fmt.Sprintf("org-%d", i+1)
|
||||
go func(ns string) {
|
||||
defer wg.Done()
|
||||
ctx := request.WithNamespace(context.Background(), ns)
|
||||
idx := cached.Get(ctx)
|
||||
require.NotNil(t, idx)
|
||||
}(namespace)
|
||||
}
|
||||
wg.Wait()
|
||||
elapsed := time.Since(startTime)
|
||||
|
||||
// Verify that all namespaces were called
|
||||
assert.Equal(t, int64(numNamespaces), provider.callCount.Load())
|
||||
|
||||
// Verify max concurrent executions shows parallelism
|
||||
maxConcurrent := provider.maxConcurrent.Load()
|
||||
assert.Equal(t, int64(numNamespaces), maxConcurrent)
|
||||
|
||||
// If all namespaces had to wait sequentially, it would take numNamespaces * blockDuration
|
||||
// With parallelism, it should be much faster (close to just blockDuration)
|
||||
sequentialTime := time.Duration(numNamespaces) * provider.blockDuration
|
||||
assert.Less(t, elapsed, sequentialTime)
|
||||
}
|
||||
|
||||
// TestCachedProvider_SameNamespaceSerialFetch verifies that the same namespace doesn't fetch concurrently
|
||||
func TestCachedProvider_SameNamespaceSerialFetch(t *testing.T) {
|
||||
// Create a blocking provider that tracks concurrent executions
|
||||
provider := &blockingProvider{
|
||||
blockDuration: 100 * time.Millisecond,
|
||||
datasources: []DataSourceInfo{
|
||||
{UID: "ds1", Type: "prometheus", Name: "Prometheus", Default: true},
|
||||
},
|
||||
}
|
||||
cached := newCachedProvider(provider.get, defaultCacheSize, time.Minute, log.New("test"))
|
||||
|
||||
numGoroutines := 10
|
||||
var wg sync.WaitGroup
|
||||
|
||||
// Launch multiple fetches for the SAME namespace simultaneously
|
||||
ctx := request.WithNamespace(context.Background(), "default")
|
||||
for i := 0; i < numGoroutines; i++ {
|
||||
wg.Add(1)
|
||||
go func() {
|
||||
defer wg.Done()
|
||||
idx := cached.Get(ctx)
|
||||
require.NotNil(t, idx)
|
||||
}()
|
||||
}
|
||||
wg.Wait()
|
||||
|
||||
// Max concurrent should be 1 since all goroutines are for the same namespace
|
||||
maxConcurrent := provider.maxConcurrent.Load()
|
||||
assert.Equal(t, int64(1), maxConcurrent)
|
||||
}
|
||||
|
||||
// blockingProvider is a test provider that simulates slow fetch operations
|
||||
// and tracks concurrent executions
|
||||
type blockingProvider struct {
|
||||
blockDuration time.Duration
|
||||
datasources []DataSourceInfo
|
||||
callCount atomic.Int64
|
||||
currentActive atomic.Int64
|
||||
maxConcurrent atomic.Int64
|
||||
}
|
||||
|
||||
func (p *blockingProvider) get(_ context.Context) any {
|
||||
p.callCount.Add(1)
|
||||
|
||||
// Track concurrent executions
|
||||
current := p.currentActive.Add(1)
|
||||
|
||||
// Update max concurrent if this is a new peak
|
||||
for {
|
||||
maxVal := p.maxConcurrent.Load()
|
||||
if current <= maxVal {
|
||||
break
|
||||
}
|
||||
if p.maxConcurrent.CompareAndSwap(maxVal, current) {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
// Simulate slow operation
|
||||
time.Sleep(p.blockDuration)
|
||||
|
||||
p.currentActive.Add(-1)
|
||||
return p.datasources
|
||||
}
|
||||
@@ -2,8 +2,9 @@ package schemaversion
|
||||
|
||||
import (
|
||||
"context"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/grafana/grafana/pkg/infra/log"
|
||||
)
|
||||
|
||||
// Shared utility functions for datasource migrations across different schema versions.
|
||||
@@ -11,65 +12,41 @@ import (
|
||||
// string names/UIDs to structured reference objects with uid, type, and apiVersion.
|
||||
|
||||
// cachedIndexProvider wraps a DataSourceIndexProvider with time-based caching.
|
||||
// This prevents multiple DB queries and index builds during operations that may call
|
||||
// provider.Index() multiple times (e.g., dashboard conversions with many datasource lookups).
|
||||
// The cache expires after 10 seconds, allowing it to be used as a long-lived singleton
|
||||
// while still refreshing periodically.
|
||||
//
|
||||
// Thread-safe: Uses sync.RWMutex to guarantee safe concurrent access.
|
||||
type cachedIndexProvider struct {
|
||||
provider DataSourceIndexProvider
|
||||
mu sync.RWMutex
|
||||
index *DatasourceIndex
|
||||
cachedAt time.Time
|
||||
cacheTTL time.Duration
|
||||
*cachedProvider[*DatasourceIndex]
|
||||
}
|
||||
|
||||
// Index returns the cached index if it's still valid (< 10s old), otherwise rebuilds it.
|
||||
// Uses RWMutex for efficient concurrent reads when cache is valid.
|
||||
// Index returns the cached index if it's still valid (< TTL old), otherwise rebuilds it.
|
||||
func (p *cachedIndexProvider) Index(ctx context.Context) *DatasourceIndex {
|
||||
// Fast path: check if cache is still valid using read lock
|
||||
p.mu.RLock()
|
||||
if p.index != nil && time.Since(p.cachedAt) < p.cacheTTL {
|
||||
idx := p.index
|
||||
p.mu.RUnlock()
|
||||
return idx
|
||||
}
|
||||
p.mu.RUnlock()
|
||||
|
||||
// Slow path: cache expired or not yet built, acquire write lock
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
// Double-check: another goroutine might have refreshed the cache
|
||||
// while we were waiting for the write lock
|
||||
if p.index != nil && time.Since(p.cachedAt) < p.cacheTTL {
|
||||
return p.index
|
||||
}
|
||||
|
||||
// Rebuild the cache
|
||||
p.index = p.provider.Index(ctx)
|
||||
p.cachedAt = time.Now()
|
||||
return p.index
|
||||
return p.Get(ctx)
|
||||
}
|
||||
|
||||
// WrapIndexProviderWithCache wraps a provider to cache the index with a 10-second TTL.
|
||||
// Useful for conversions or migrations that may call provider.Index() multiple times.
|
||||
// The cache expires after 10 seconds, making it suitable for use as a long-lived singleton
|
||||
// at the top level of dependency injection while still refreshing periodically.
|
||||
//
|
||||
// Example usage in dashboard conversion:
|
||||
//
|
||||
// cachedDsIndexProvider := schemaversion.WrapIndexProviderWithCache(dsIndexProvider)
|
||||
// // Now all calls to cachedDsIndexProvider.Index(ctx) return the same cached index
|
||||
// // for up to 10 seconds before refreshing
|
||||
func WrapIndexProviderWithCache(provider DataSourceIndexProvider) DataSourceIndexProvider {
|
||||
if provider == nil {
|
||||
return nil
|
||||
// cachedLibraryElementProvider wraps a LibraryElementIndexProvider with time-based caching.
|
||||
type cachedLibraryElementProvider struct {
|
||||
*cachedProvider[[]LibraryElementInfo]
|
||||
}
|
||||
|
||||
func (p *cachedLibraryElementProvider) GetLibraryElementInfo(ctx context.Context) []LibraryElementInfo {
|
||||
return p.Get(ctx)
|
||||
}
|
||||
|
||||
// WrapIndexProviderWithCache wraps a DataSourceIndexProvider to cache indexes with a configurable TTL.
|
||||
func WrapIndexProviderWithCache(provider DataSourceIndexProvider, cacheTTL time.Duration) DataSourceIndexProvider {
|
||||
if provider == nil || cacheTTL <= 0 {
|
||||
return provider
|
||||
}
|
||||
return &cachedIndexProvider{
|
||||
provider: provider,
|
||||
cacheTTL: 10 * time.Second,
|
||||
newCachedProvider[*DatasourceIndex](provider.Index, defaultCacheSize, cacheTTL, log.New("schemaversion.dsindexprovider")),
|
||||
}
|
||||
}
|
||||
|
||||
// WrapLibraryElementProviderWithCache wraps a LibraryElementIndexProvider to cache library elements with a configurable TTL.
|
||||
func WrapLibraryElementProviderWithCache(provider LibraryElementIndexProvider, cacheTTL time.Duration) LibraryElementIndexProvider {
|
||||
if provider == nil || cacheTTL <= 0 {
|
||||
return provider
|
||||
}
|
||||
return &cachedLibraryElementProvider{
|
||||
newCachedProvider[[]LibraryElementInfo](provider.GetLibraryElementInfo, defaultCacheSize, cacheTTL, log.New("schemaversion.leindexprovider")),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -216,60 +193,3 @@ func MigrateDatasourceNameToRef(nameOrRef interface{}, options map[string]bool,
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
// cachedLibraryElementProvider wraps a LibraryElementIndexProvider with time-based caching.
|
||||
// This prevents multiple DB queries during operations that may call GetLibraryElementInfo()
|
||||
// multiple times (e.g., dashboard conversions with many library panel lookups).
|
||||
// The cache expires after 10 seconds, allowing it to be used as a long-lived singleton
|
||||
// while still refreshing periodically.
|
||||
//
|
||||
// Thread-safe: Uses sync.RWMutex to guarantee safe concurrent access.
|
||||
type cachedLibraryElementProvider struct {
|
||||
provider LibraryElementIndexProvider
|
||||
mu sync.RWMutex
|
||||
elements []LibraryElementInfo
|
||||
cachedAt time.Time
|
||||
cacheTTL time.Duration
|
||||
}
|
||||
|
||||
// GetLibraryElementInfo returns the cached library elements if they're still valid (< 10s old), otherwise rebuilds the cache.
|
||||
// Uses RWMutex for efficient concurrent reads when cache is valid.
|
||||
func (p *cachedLibraryElementProvider) GetLibraryElementInfo(ctx context.Context) []LibraryElementInfo {
|
||||
// Fast path: check if cache is still valid using read lock
|
||||
p.mu.RLock()
|
||||
if p.elements != nil && time.Since(p.cachedAt) < p.cacheTTL {
|
||||
elements := p.elements
|
||||
p.mu.RUnlock()
|
||||
return elements
|
||||
}
|
||||
p.mu.RUnlock()
|
||||
|
||||
// Slow path: cache expired or not yet built, acquire write lock
|
||||
p.mu.Lock()
|
||||
defer p.mu.Unlock()
|
||||
|
||||
// Double-check: another goroutine might have refreshed the cache
|
||||
// while we were waiting for the write lock
|
||||
if p.elements != nil && time.Since(p.cachedAt) < p.cacheTTL {
|
||||
return p.elements
|
||||
}
|
||||
|
||||
// Rebuild the cache
|
||||
p.elements = p.provider.GetLibraryElementInfo(ctx)
|
||||
p.cachedAt = time.Now()
|
||||
return p.elements
|
||||
}
|
||||
|
||||
// WrapLibraryElementProviderWithCache wraps a provider to cache library elements with a 10-second TTL.
|
||||
// Useful for conversions or migrations that may call GetLibraryElementInfo() multiple times.
|
||||
// The cache expires after 10 seconds, making it suitable for use as a long-lived singleton
|
||||
// at the top level of dependency injection while still refreshing periodically.
|
||||
func WrapLibraryElementProviderWithCache(provider LibraryElementIndexProvider) LibraryElementIndexProvider {
|
||||
if provider == nil {
|
||||
return nil
|
||||
}
|
||||
return &cachedLibraryElementProvider{
|
||||
provider: provider,
|
||||
cacheTTL: 10 * time.Second,
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user