diff --git a/packages/grafana-data/src/types/featureToggles.gen.ts b/packages/grafana-data/src/types/featureToggles.gen.ts index 3b9875c801e..9582065043f 100644 --- a/packages/grafana-data/src/types/featureToggles.gen.ts +++ b/packages/grafana-data/src/types/featureToggles.gen.ts @@ -353,6 +353,10 @@ export interface FeatureToggles { */ alertmanagerRemoteSecondary?: boolean; /** + * Enables a feature to avoid issues with concurrent writes to the alerting provenance table in MySQL + */ + alertingProvenanceLockWrites?: boolean; + /** * Enable Grafana to have a remote Alertmanager instance as the primary Alertmanager. */ alertmanagerRemotePrimary?: boolean; diff --git a/pkg/services/featuremgmt/registry.go b/pkg/services/featuremgmt/registry.go index d91e76998e2..a5881618107 100644 --- a/pkg/services/featuremgmt/registry.go +++ b/pkg/services/featuremgmt/registry.go @@ -585,6 +585,14 @@ var ( Stage: FeatureStageExperimental, Owner: grafanaAlertingSquad, }, + { + Name: "alertingProvenanceLockWrites", + Description: "Enables a feature to avoid issues with concurrent writes to the alerting provenance table in MySQL", + Stage: FeatureStageExperimental, + Owner: grafanaAlertingSquad, + HideFromAdminPage: true, + HideFromDocs: true, + }, { Name: "alertmanagerRemotePrimary", Description: "Enable Grafana to have a remote Alertmanager instance as the primary Alertmanager.", diff --git a/pkg/services/featuremgmt/toggles_gen.csv b/pkg/services/featuremgmt/toggles_gen.csv index 689b0883024..8bf9cca2706 100644 --- a/pkg/services/featuremgmt/toggles_gen.csv +++ b/pkg/services/featuremgmt/toggles_gen.csv @@ -77,6 +77,7 @@ cachingOptimizeSerializationMemoryUsage,experimental,@grafana/grafana-operator-e prometheusCodeModeMetricNamesSearch,experimental,@grafana/oss-big-tent,false,false,true addFieldFromCalculationStatFunctions,GA,@grafana/dataviz-squad,false,false,true alertmanagerRemoteSecondary,experimental,@grafana/alerting-squad,false,false,false +alertingProvenanceLockWrites,experimental,@grafana/alerting-squad,false,false,false alertmanagerRemotePrimary,experimental,@grafana/alerting-squad,false,false,false annotationPermissionUpdate,GA,@grafana/identity-access-team,false,false,false extractFieldsNameDeduplication,experimental,@grafana/dataviz-squad,false,false,true diff --git a/pkg/services/featuremgmt/toggles_gen.go b/pkg/services/featuremgmt/toggles_gen.go index 5ea9f0f1995..5e8ccbb5615 100644 --- a/pkg/services/featuremgmt/toggles_gen.go +++ b/pkg/services/featuremgmt/toggles_gen.go @@ -319,6 +319,10 @@ const ( // Enable Grafana to sync configuration and state with a remote Alertmanager. FlagAlertmanagerRemoteSecondary = "alertmanagerRemoteSecondary" + // FlagAlertingProvenanceLockWrites + // Enables a feature to avoid issues with concurrent writes to the alerting provenance table in MySQL + FlagAlertingProvenanceLockWrites = "alertingProvenanceLockWrites" + // FlagAlertmanagerRemotePrimary // Enable Grafana to have a remote Alertmanager instance as the primary Alertmanager. FlagAlertmanagerRemotePrimary = "alertmanagerRemotePrimary" diff --git a/pkg/services/featuremgmt/toggles_gen.json b/pkg/services/featuremgmt/toggles_gen.json index 2b3a1c27cf9..88ab6c38de0 100644 --- a/pkg/services/featuremgmt/toggles_gen.json +++ b/pkg/services/featuremgmt/toggles_gen.json @@ -351,6 +351,20 @@ "frontend": true } }, + { + "metadata": { + "name": "alertingProvenanceLockWrites", + "resourceVersion": "1753284360846", + "creationTimestamp": "2025-07-23T15:26:00Z" + }, + "spec": { + "description": "Enables a feature to avoid issues with concurrent writes to the alerting provenance table in MySQL", + "stage": "experimental", + "codeowner": "@grafana/alerting-squad", + "hideFromAdminPage": true, + "hideFromDocs": true + } + }, { "metadata": { "name": "alertingQueryAndExpressionsStepMode", diff --git a/pkg/services/ngalert/store/provisioning_store.go b/pkg/services/ngalert/store/provisioning_store.go index d6999eec31d..f1306168134 100644 --- a/pkg/services/ngalert/store/provisioning_store.go +++ b/pkg/services/ngalert/store/provisioning_store.go @@ -5,6 +5,7 @@ import ( "fmt" "github.com/grafana/grafana/pkg/infra/db" + "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/ngalert/models" ) @@ -70,25 +71,61 @@ func (st DBstore) SetProvenance(ctx context.Context, o models.Provisionable, org // TODO: Add a unit-of-work pattern, so updating objects + provenance will happen consistently with rollbacks across stores. // TODO: Need to make sure that writing a record where our concurrency key fails will also fail the whole transaction. That way, this gets rolled back too. can't just check that 0 updates happened inmemory. Check with jp. If not possible, we need our own concurrency key. // TODO: Clean up stale provenance records periodically. - upsertSQL := st.SQLStore.GetDialect().UpsertSQL( - provenanceRecord{}.TableName(), - []string{"record_key", "record_type", "org_id"}, - []string{"record_key", "record_type", "org_id", "provenance"}) - params := []interface{}{ - recordKey, - recordType, - org, - p, + if st.FeatureToggles.IsEnabledGlobally(featuremgmt.FlagAlertingProvenanceLockWrites) { + return st.setProvenanceWithLocking(sess, recordKey, recordType, org, p) } + return st.setProvenanceUpsert(sess, recordKey, recordType, org, p) + }) +} - _, err := sess.SQL(upsertSQL, params...).Query() +func (st DBstore) setProvenanceUpsert(sess *db.Session, recordKey, recordType string, org int64, p models.Provenance) error { + upsertSQL := st.SQLStore.GetDialect().UpsertSQL( + provenanceRecord{}.TableName(), + []string{"record_key", "record_type", "org_id"}, + []string{"record_key", "record_type", "org_id", "provenance"}) + + params := []interface{}{ + recordKey, + recordType, + org, + p, + } + + _, err := sess.SQL(upsertSQL, params...).Query() + if err != nil { + return fmt.Errorf("failed to store provisioning status: %w", err) + } + return nil +} + +func (st DBstore) setProvenanceWithLocking(sess *db.Session, recordKey, recordType string, org int64, p models.Provenance) error { + // Check if the record exists with FOR UPDATE lock. + // If it does, we just update, otherwise we upsert the record. + // This is done to avoid deadlocks that can occur in MySQL when multiple transactions try to + // insert records (even different) because of the gap and insert intention locks. + exists, err := sess.Table(provenanceRecord{}). + Where("record_key = ? AND record_type = ? AND org_id = ?", recordKey, recordType, org). + ForUpdate(). + Exist() + if err != nil { + return fmt.Errorf("failed to check if provenance record exists: %w", err) + } + + if exists { + _, err = sess.Table(provenanceRecord{}). + Where("record_key = ? AND record_type = ? AND org_id = ?", recordKey, recordType, org). + Update(map[string]interface{}{ + "provenance": p, + }) if err != nil { return fmt.Errorf("failed to store provisioning status: %w", err) } - return nil - }) + } + + // Still upsert in case it was created while we were checking + return st.setProvenanceUpsert(sess, recordKey, recordType, org, p) } // DeleteProvenance deletes the provenance record from the table diff --git a/pkg/services/ngalert/store/provisioning_store_test.go b/pkg/services/ngalert/store/provisioning_store_test.go index 568f870027d..ea4daf9811b 100644 --- a/pkg/services/ngalert/store/provisioning_store_test.go +++ b/pkg/services/ngalert/store/provisioning_store_test.go @@ -2,10 +2,15 @@ package store_test import ( "context" + "fmt" + "math/rand" + "sync" "testing" + "time" "github.com/stretchr/testify/require" + "github.com/grafana/grafana/pkg/services/featuremgmt" "github.com/grafana/grafana/pkg/services/ngalert" "github.com/grafana/grafana/pkg/services/ngalert/models" "github.com/grafana/grafana/pkg/services/ngalert/provisioning" @@ -19,113 +24,206 @@ func TestIntegrationProvisioningStore(t *testing.T) { if testing.Short() { t.Skip("skipping integration test") } - store := createProvisioningStoreSut(tests.SetupTestEnv(t, testAlertingIntervalSeconds)) - t.Run("Default provenance of a known type is None", func(t *testing.T) { - rule := models.AlertRule{ - UID: "asdf", + testCases := []struct { + name string + featureEnabled bool + }{ + { + name: "without feature flag", + featureEnabled: false, + }, + { + name: "with feature flag", + featureEnabled: true, + }, + } + + for _, tc := range testCases { + t.Run(tc.name, func(t *testing.T) { + ng, dbStore := tests.SetupTestEnv(t, testAlertingIntervalSeconds) + if tc.featureEnabled { + dbStore.FeatureToggles = featuremgmt.WithFeatures(featuremgmt.FlagAlertingProvenanceLockWrites) + } + store := createProvisioningStoreSut(ng, dbStore) + + t.Run("Default provenance of a known type is None", func(t *testing.T) { + rule := models.AlertRule{ + UID: "asdf", + } + + provenance, err := store.GetProvenance(context.Background(), &rule, 1) + + require.NoError(t, err) + require.Equal(t, models.ProvenanceNone, provenance) + }) + + t.Run("Store returns saved provenance type", func(t *testing.T) { + rule := models.AlertRule{ + UID: "123", + } + err := store.SetProvenance(context.Background(), &rule, 1, models.ProvenanceFile) + require.NoError(t, err) + + p, err := store.GetProvenance(context.Background(), &rule, 1) + + require.NoError(t, err) + require.Equal(t, models.ProvenanceFile, p) + }) + + t.Run("Store does not get provenance of record with different org ID", func(t *testing.T) { + ruleOrg2 := models.AlertRule{ + UID: "456", + } + ruleOrg3 := models.AlertRule{ + UID: "456", + } + err := store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceFile) + require.NoError(t, err) + + p, err := store.GetProvenance(context.Background(), &ruleOrg3, 3) + + require.NoError(t, err) + require.Equal(t, models.ProvenanceNone, p) + }) + + t.Run("Store only updates provenance of record with given org ID", func(t *testing.T) { + ruleOrg2 := models.AlertRule{ + UID: "789", + OrgID: 2, + } + ruleOrg3 := models.AlertRule{ + UID: "789", + OrgID: 3, + } + err := store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceFile) + require.NoError(t, err) + err = store.SetProvenance(context.Background(), &ruleOrg3, 3, models.ProvenanceFile) + require.NoError(t, err) + + err = store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceAPI) + require.NoError(t, err) + + p, err := store.GetProvenance(context.Background(), &ruleOrg2, 2) + require.NoError(t, err) + require.Equal(t, models.ProvenanceAPI, p) + p, err = store.GetProvenance(context.Background(), &ruleOrg3, 3) + require.NoError(t, err) + require.Equal(t, models.ProvenanceFile, p) + }) + + t.Run("Store should return all provenances by type", func(t *testing.T) { + const orgID = 123 + rule1 := models.AlertRule{ + UID: "789", + OrgID: orgID, + } + rule2 := models.AlertRule{ + UID: "790", + OrgID: orgID, + } + err := store.SetProvenance(context.Background(), &rule1, orgID, models.ProvenanceFile) + require.NoError(t, err) + err = store.SetProvenance(context.Background(), &rule2, orgID, models.ProvenanceAPI) + require.NoError(t, err) + + p, err := store.GetProvenances(context.Background(), orgID, rule1.ResourceType()) + require.NoError(t, err) + require.Len(t, p, 2) + require.Equal(t, models.ProvenanceFile, p[rule1.UID]) + require.Equal(t, models.ProvenanceAPI, p[rule2.UID]) + }) + + t.Run("Store should delete provenance correctly", func(t *testing.T) { + const orgID = 1234 + ruleOrg := models.AlertRule{ + UID: "7834539", + OrgID: orgID, + } + err := store.SetProvenance(context.Background(), &ruleOrg, orgID, models.ProvenanceFile) + require.NoError(t, err) + p, err := store.GetProvenance(context.Background(), &ruleOrg, orgID) + require.NoError(t, err) + require.Equal(t, models.ProvenanceFile, p) + + err = store.DeleteProvenance(context.Background(), &ruleOrg, orgID) + require.NoError(t, err) + + p, err = store.GetProvenance(context.Background(), &ruleOrg, orgID) + require.NoError(t, err) + require.Equal(t, models.ProvenanceNone, p) + }) + }) + } +} + +func TestSetProvenance_DeadlockScenarios(t *testing.T) { + if testing.Short() { + t.Skip("skipping integration test") + } + + ng, dbStore := tests.SetupTestEnv(t, testAlertingIntervalSeconds) + dbStore.FeatureToggles = featuremgmt.WithFeatures(featuremgmt.FlagAlertingProvenanceLockWrites) + store := createProvisioningStoreSut(ng, dbStore) + concurrency := 20 + + t.Run("Same record, different orgs", func(t *testing.T) { + rule := &models.AlertRule{UID: "same-record-diff-orgs"} + + var wg sync.WaitGroup + for i := range concurrency { + wg.Add(1) + go func(orgID int64) { + defer wg.Done() + time.Sleep(time.Microsecond * time.Duration(rand.Intn(100))) + err := store.SetProvenance(context.Background(), rule, orgID, models.ProvenanceAPI) + require.NoError(t, err) + }(int64(i + 1)) } - - provenance, err := store.GetProvenance(context.Background(), &rule, 1) - - require.NoError(t, err) - require.Equal(t, models.ProvenanceNone, provenance) + wg.Wait() }) - t.Run("Store returns saved provenance type", func(t *testing.T) { - rule := models.AlertRule{ - UID: "123", + t.Run("Different records, same org", func(t *testing.T) { + orgID := int64(1) + + var wg sync.WaitGroup + for i := range concurrency { + wg.Add(1) + go func(uid string) { + defer wg.Done() + time.Sleep(time.Microsecond * time.Duration(rand.Intn(100))) + rule := &models.AlertRule{UID: uid} + err := store.SetProvenance(context.Background(), rule, orgID, models.ProvenanceFile) + require.NoError(t, err) + }(fmt.Sprintf("diff-record-same-org-%d", i)) } - err := store.SetProvenance(context.Background(), &rule, 1, models.ProvenanceFile) - require.NoError(t, err) - - p, err := store.GetProvenance(context.Background(), &rule, 1) - - require.NoError(t, err) - require.Equal(t, models.ProvenanceFile, p) + wg.Wait() }) - t.Run("Store does not get provenance of record with different org ID", func(t *testing.T) { - ruleOrg2 := models.AlertRule{ - UID: "456", + t.Run("Mixed operations", func(t *testing.T) { + rule := &models.AlertRule{UID: "mixed-ops"} + orgID := int64(1) + + var wg sync.WaitGroup + // Mix SetProvenance and GetProvenance operations + for range concurrency { + wg.Add(1) + go func() { + defer wg.Done() + time.Sleep(time.Microsecond * time.Duration(rand.Intn(100))) + err := store.SetProvenance(context.Background(), rule, orgID, models.ProvenanceAPI) + require.NoError(t, err) + }() + + wg.Add(1) + go func() { + defer wg.Done() + time.Sleep(time.Microsecond * time.Duration(rand.Intn(100))) + _, err := store.GetProvenance(context.Background(), rule, orgID) + require.NoError(t, err) + }() } - ruleOrg3 := models.AlertRule{ - UID: "456", - } - err := store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceFile) - require.NoError(t, err) - - p, err := store.GetProvenance(context.Background(), &ruleOrg3, 3) - - require.NoError(t, err) - require.Equal(t, models.ProvenanceNone, p) - }) - - t.Run("Store only updates provenance of record with given org ID", func(t *testing.T) { - ruleOrg2 := models.AlertRule{ - UID: "789", - OrgID: 2, - } - ruleOrg3 := models.AlertRule{ - UID: "789", - OrgID: 3, - } - err := store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceFile) - require.NoError(t, err) - err = store.SetProvenance(context.Background(), &ruleOrg3, 3, models.ProvenanceFile) - require.NoError(t, err) - - err = store.SetProvenance(context.Background(), &ruleOrg2, 2, models.ProvenanceAPI) - require.NoError(t, err) - - p, err := store.GetProvenance(context.Background(), &ruleOrg2, 2) - require.NoError(t, err) - require.Equal(t, models.ProvenanceAPI, p) - p, err = store.GetProvenance(context.Background(), &ruleOrg3, 3) - require.NoError(t, err) - require.Equal(t, models.ProvenanceFile, p) - }) - - t.Run("Store should return all provenances by type", func(t *testing.T) { - const orgID = 123 - rule1 := models.AlertRule{ - UID: "789", - OrgID: orgID, - } - rule2 := models.AlertRule{ - UID: "790", - OrgID: orgID, - } - err := store.SetProvenance(context.Background(), &rule1, orgID, models.ProvenanceFile) - require.NoError(t, err) - err = store.SetProvenance(context.Background(), &rule2, orgID, models.ProvenanceAPI) - require.NoError(t, err) - - p, err := store.GetProvenances(context.Background(), orgID, rule1.ResourceType()) - require.NoError(t, err) - require.Len(t, p, 2) - require.Equal(t, models.ProvenanceFile, p[rule1.UID]) - require.Equal(t, models.ProvenanceAPI, p[rule2.UID]) - }) - - t.Run("Store should delete provenance correctly", func(t *testing.T) { - const orgID = 1234 - ruleOrg := models.AlertRule{ - UID: "7834539", - OrgID: orgID, - } - err := store.SetProvenance(context.Background(), &ruleOrg, orgID, models.ProvenanceFile) - require.NoError(t, err) - p, err := store.GetProvenance(context.Background(), &ruleOrg, orgID) - require.NoError(t, err) - require.Equal(t, models.ProvenanceFile, p) - - err = store.DeleteProvenance(context.Background(), &ruleOrg, orgID) - require.NoError(t, err) - - p, err = store.GetProvenance(context.Background(), &ruleOrg, orgID) - require.NoError(t, err) - require.Equal(t, models.ProvenanceNone, p) + wg.Wait() }) }