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
This commit is contained in:
@@ -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
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user