From 0eb28d4f370520561d515ad895a92859550a1a7d Mon Sep 17 00:00:00 2001 From: Gabriel MABILLE Date: Thu, 16 Oct 2025 08:33:06 +0200 Subject: [PATCH] `grafana-iam`: Async write to zanzana (#112357) * `grafana-iam`: Async write to zanzana * More succint * Add a metric to keep track of waiting times for future calibration * metrics --- pkg/registry/apis/iam/hooks.go | 91 ++++++++++++++++------------- pkg/registry/apis/iam/hooks_test.go | 10 +++- pkg/registry/apis/iam/metrics.go | 32 ++++++++++ pkg/registry/apis/iam/models.go | 2 + pkg/registry/apis/iam/register.go | 11 ++-- 5 files changed, 99 insertions(+), 47 deletions(-) create mode 100644 pkg/registry/apis/iam/metrics.go diff --git a/pkg/registry/apis/iam/hooks.go b/pkg/registry/apis/iam/hooks.go index e8bdbc33515..f572711512b 100644 --- a/pkg/registry/apis/iam/hooks.go +++ b/pkg/registry/apis/iam/hooks.go @@ -100,58 +100,71 @@ func (b *IdentityAccessManagementAPIBuilder) AfterResourcePermissionCreate(obj r return } + // Grab a ticket to write to Zanzana + // This limits the amount of concurrent writes to Zanzana + wait := time.Now() + b.zTickets <- true + hooksWaitHistogram.Observe(time.Since(wait).Seconds()) // Record wait time + rp, ok := obj.(*iamv0.ResourcePermission) if !ok { return } - resource := rp.Spec.Resource - permissions := rp.Spec.Permissions + go func(rp *iamv0.ResourcePermission) { + defer func() { + // Release the ticket after write is done + <-b.zTickets + }() - object := zanzana.NewObjectEntry(toZanzanaType(resource.ApiGroup), resource.ApiGroup, resource.Resource, "", resource.Name) + resource := rp.Spec.Resource + permissions := rp.Spec.Permissions - tuples := make([]*v1.TupleKey, 0, len(permissions)) - for _, p := range permissions { - tuple, err := NewResourceTuple(object, resource, p) - if err != nil { - b.logger.Error("failed to create resource permission tuple", - "namespace", rp.Namespace, - "object", object, - "err", err, - ) + object := zanzana.NewObjectEntry(toZanzanaType(resource.ApiGroup), resource.ApiGroup, resource.Resource, "", resource.Name) - continue + tuples := make([]*v1.TupleKey, 0, len(permissions)) + for _, p := range permissions { + tuple, err := NewResourceTuple(object, resource, p) + if err != nil { + b.logger.Error("failed to create resource permission tuple", + "namespace", rp.Namespace, + "object", object, + "err", err, + ) + + continue + } + tuples = append(tuples, tuple) } - tuples = append(tuples, tuple) - } - // Avoid writing if there are no valid tuples - if len(tuples) == 0 { - b.logger.Warn("no valid tuples to write", "namespace", rp.Namespace, "resource", object) - return - } + // Avoid writing if there are no valid tuples + if len(tuples) == 0 { + b.logger.Warn("no valid tuples to write", "namespace", rp.Namespace, "resource", object) + return + } - b.logger.Debug("writing resource permission to zanzana", - "namespace", rp.Namespace, - "object", object, - "tuplesCnt", len(tuples), - ) - - ctx, cancel := context.WithTimeout(context.Background(), defaultWriteTimeout) - defer cancel() - - err := b.zClient.Write(ctx, &v1.WriteRequest{ - Namespace: rp.Namespace, - Writes: &v1.WriteRequestWrites{ - TupleKeys: tuples, - }, - }) - if err != nil { - b.logger.Error("failed to write resource permission to zanzana", - "err", err, + b.logger.Debug("writing resource permission to zanzana", "namespace", rp.Namespace, "object", object, "tuplesCnt", len(tuples), ) - } + + ctx, cancel := context.WithTimeout(context.Background(), defaultWriteTimeout) + defer cancel() + + err := b.zClient.Write(ctx, &v1.WriteRequest{ + Namespace: rp.Namespace, + Writes: &v1.WriteRequestWrites{ + TupleKeys: tuples, + }, + }) + if err != nil { + b.logger.Error("failed to write resource permission to zanzana", + "err", err, + "namespace", rp.Namespace, + "object", object, + "tuplesCnt", len(tuples), + ) + } + }(rp.DeepCopy()) // Pass a copy of the object } diff --git a/pkg/registry/apis/iam/hooks_test.go b/pkg/registry/apis/iam/hooks_test.go index 8c11f2faa8b..3ee0c7bcd24 100644 --- a/pkg/registry/apis/iam/hooks_test.go +++ b/pkg/registry/apis/iam/hooks_test.go @@ -25,7 +25,8 @@ func (f *FakeZanzanaClient) Write(ctx context.Context, req *v1.WriteRequest) err func TestAfterResourcePermissionCreate(t *testing.T) { b := &IdentityAccessManagementAPIBuilder{ - logger: log.NewNopLogger(), + logger: log.NewNopLogger(), + zTickets: make(chan bool, 1), } t.Run("should create zanzana entries for folder resource permissions", func(t *testing.T) { folderPerm := iamv0.ResourcePermission{ @@ -47,7 +48,7 @@ func TestAfterResourcePermissionCreate(t *testing.T) { require.NotNil(t, req) require.NotNil(t, req.Writes) require.Len(t, req.Writes.TupleKeys, 2) - require.Equal(t, req.Namespace, "org-2") + require.Equal(t, "org-2", req.Namespace) require.Equal( t, req.Writes.TupleKeys[0], @@ -65,6 +66,9 @@ func TestAfterResourcePermissionCreate(t *testing.T) { b.AfterResourcePermissionCreate(&folderPerm, nil) }) + // Wait for the ticket to be released + <-b.zTickets + t.Run("should create zanzana entries for dashboard resource permissions", func(t *testing.T) { dashPerm := iamv0.ResourcePermission{ ObjectMeta: metav1.ObjectMeta{ @@ -87,7 +91,7 @@ func TestAfterResourcePermissionCreate(t *testing.T) { require.NotNil(t, req) require.NotNil(t, req.Writes) require.Len(t, req.Writes.TupleKeys, 2) - require.Equal(t, req.Namespace, "default") + require.Equal(t, "default", req.Namespace) tuple1 := req.Writes.TupleKeys[0] require.NotNil(t, tuple1.Condition) diff --git a/pkg/registry/apis/iam/metrics.go b/pkg/registry/apis/iam/metrics.go new file mode 100644 index 00000000000..6c25a8746f7 --- /dev/null +++ b/pkg/registry/apis/iam/metrics.go @@ -0,0 +1,32 @@ +package iam + +import ( + "sync" + + "github.com/grafana/grafana/pkg/infra/log" + "github.com/prometheus/client_golang/prometheus" +) + +const ( + metricsNamespace = "iam" + metricsSubSystem = "apiserver" +) + +var ( + registerOnce sync.Once + hooksWaitHistogram = prometheus.NewHistogram(prometheus.HistogramOpts{ + Namespace: metricsNamespace, + Subsystem: metricsSubSystem, + Name: "hooks_wait_duration_seconds", + Help: "Time spent in the hooks waiting for a ticket to start processing", + Buckets: prometheus.ExponentialBuckets(0.001, 2, 5), // 1ms to ~16s + }) +) + +func registerMetrics(reg prometheus.Registerer) { + registerOnce.Do(func() { + if err := reg.Register(hooksWaitHistogram); err != nil { + log.New("iam.apis").Warn("failed to register iam apiserver metrics", "error", err) + } + }) +} diff --git a/pkg/registry/apis/iam/models.go b/pkg/registry/apis/iam/models.go index 49e1662263a..0ab7869c30a 100644 --- a/pkg/registry/apis/iam/models.go +++ b/pkg/registry/apis/iam/models.go @@ -53,6 +53,8 @@ type IdentityAccessManagementAPIBuilder struct { // - permissions // - assignments zClient zanzana.Client + // Buffered channel to limit the amount of concurrent writes to Zanzana + zTickets chan bool reg prometheus.Registerer logger log.Logger diff --git a/pkg/registry/apis/iam/register.go b/pkg/registry/apis/iam/register.go index e5c1d2b98f2..9ab61ff8a1c 100644 --- a/pkg/registry/apis/iam/register.go +++ b/pkg/registry/apis/iam/register.go @@ -46,6 +46,8 @@ import ( "github.com/grafana/grafana/pkg/storage/unified/resource" ) +const MaxConcurrentZanzanaWrites = 20 + func RegisterAPIService( features featuremgmt.FeatureToggles, apiregistration builder.APIRegistrar, @@ -63,6 +65,7 @@ func RegisterAPIService( store := legacy.NewLegacySQLStores(dbProvider) legacyAccessClient := newLegacyAccessClient(ac, store) authorizer := newIAMAuthorizer(accessClient, legacyAccessClient) + registerMetrics(reg) builder := &IdentityAccessManagementAPIBuilder{ store: store, @@ -75,21 +78,19 @@ func RegisterAPIService( legacyAccessClient: legacyAccessClient, accessClient: accessClient, zClient: zClient, + zTickets: make(chan bool, MaxConcurrentZanzanaWrites), display: user.NewLegacyDisplayREST(store), reg: reg, logger: log.New("iam.apis"), features: features, - // enableAuthZApis: features.IsEnabledGlobally(featuremgmt.FlagKubernetesAuthzApis), - // enableResourcePermissionApis: features.IsEnabledGlobally(featuremgmt.FlagKubernetesAuthzResourcePermissionApis), - // enableAuthnMutation: features.IsEnabledGlobally(featuremgmt.FlagKubernetesAuthnMutation), - enableDualWriter: true, + enableDualWriter: true, } apiregistration.RegisterAPI(builder) return builder, nil } -// TODO zClient, reg +// TODO zClient, zTickets, reg func NewAPIService( accessClient types.AccessClient, dbProvider legacysql.LegacyDatabaseProvider,