Alerting: Fix deadlocks in provenance table (#106370)
What is this feature? This PR further improves concurrent updates to the provenance table (a follow-up for #101688). The fix explicitly checks for a provenance record using SELECT ... FOR UPDATE for the databases that support it before performing either an update or upsert operation. This preemptive locking reduces the possibility of deadlocks in MySQL. This works only with the alertingProvenanceLockWrites feature flag enabled. Why do we need this feature? The current implementation (directly performing upserts without prior locking) may still encounter deadlocks because of the gap and insert-intention locks in some configurations, for example, with repeatable-read transaction isolation level in MySQL+InnoDB. --------- Co-authored-by: William Wernert <william.wernert@grafana.com>
This commit is contained in:
co-authored by
William Wernert
parent
709aeb4e1a
commit
7c43c061a8
@@ -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;
|
||||
|
||||
@@ -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.",
|
||||
|
||||
@@ -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
|
||||
|
||||
|
@@ -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"
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user