* Alerting: Fetch alert rule provenances for a page of rules only * error when failed to fetch provenance
167 lines
5.9 KiB
Go
167 lines
5.9 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
|
|
"github.com/grafana/grafana/pkg/infra/db"
|
|
"github.com/grafana/grafana/pkg/services/featuremgmt"
|
|
"github.com/grafana/grafana/pkg/services/ngalert/models"
|
|
)
|
|
|
|
type provenanceRecord struct {
|
|
Id int `xorm:"pk autoincr 'id'"`
|
|
OrgID int64 `xorm:"'org_id'"`
|
|
RecordKey string
|
|
RecordType string
|
|
Provenance models.Provenance
|
|
}
|
|
|
|
func (pr provenanceRecord) TableName() string {
|
|
return "provenance_type"
|
|
}
|
|
|
|
// GetProvenance gets the provenance status for a provisionable object.
|
|
func (st DBstore) GetProvenance(ctx context.Context, o models.Provisionable, org int64) (models.Provenance, error) {
|
|
recordType := o.ResourceType()
|
|
recordKey := o.ResourceID()
|
|
|
|
provenance := models.ProvenanceNone
|
|
err := st.SQLStore.WithDbSession(ctx, func(sess *db.Session) error {
|
|
filter := "record_key = ? AND record_type = ? AND org_id = ?"
|
|
var result models.Provenance
|
|
has, err := sess.Table(provenanceRecord{}).Where(filter, recordKey, recordType, org).Desc("id").Cols("provenance").Get(&result)
|
|
if err != nil {
|
|
return fmt.Errorf("failed to query for existing provenance status: %w", err)
|
|
}
|
|
if has {
|
|
provenance = result
|
|
}
|
|
return nil
|
|
})
|
|
if err != nil {
|
|
return models.ProvenanceNone, err
|
|
}
|
|
return provenance, nil
|
|
}
|
|
|
|
// GetProvenance gets the provenance status for a provisionable object.
|
|
func (st DBstore) GetProvenances(ctx context.Context, org int64, resourceType string) (map[string]models.Provenance, error) {
|
|
resultMap := make(map[string]models.Provenance)
|
|
err := st.SQLStore.WithDbSession(ctx, func(sess *db.Session) error {
|
|
filter := "record_type = ? AND org_id = ?"
|
|
rawData, err := sess.Table(provenanceRecord{}).Where(filter, resourceType, org).Desc("id").Cols("record_key", "provenance").QueryString()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to query for existing provenance status: %w", err)
|
|
}
|
|
for _, data := range rawData {
|
|
resultMap[data["record_key"]] = models.Provenance(data["provenance"])
|
|
}
|
|
return nil
|
|
})
|
|
return resultMap, err
|
|
}
|
|
|
|
// GetProvenancesByUIDs gets the provenance status for specific UIDs.
|
|
func (st DBstore) GetProvenancesByUIDs(ctx context.Context, org int64, resourceType string, uids []string) (map[string]models.Provenance, error) {
|
|
if len(uids) == 0 {
|
|
return map[string]models.Provenance{}, nil
|
|
}
|
|
|
|
result := make(map[string]models.Provenance, len(uids))
|
|
err := st.SQLStore.WithDbSession(ctx, func(sess *db.Session) error {
|
|
rawData, err := sess.Table(provenanceRecord{}).
|
|
Where("record_type = ? AND org_id = ?", resourceType, org).
|
|
In("record_key", uids).
|
|
Cols("record_key", "provenance").
|
|
QueryString()
|
|
if err != nil {
|
|
return fmt.Errorf("failed to query for existing provenance status: %w", err)
|
|
}
|
|
for _, data := range rawData {
|
|
result[data["record_key"]] = models.Provenance(data["provenance"])
|
|
}
|
|
return nil
|
|
})
|
|
return result, err
|
|
}
|
|
|
|
// SetProvenance changes the provenance status for a provisionable object.
|
|
func (st DBstore) SetProvenance(ctx context.Context, o models.Provisionable, org int64, p models.Provenance) error {
|
|
recordType := o.ResourceType()
|
|
recordKey := o.ResourceID()
|
|
|
|
return st.SQLStore.WithTransactionalDbSession(ctx, func(sess *db.Session) error {
|
|
// 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.
|
|
|
|
//nolint:staticcheck // not yet migrated to OpenFeature
|
|
if st.FeatureToggles.IsEnabledGlobally(featuremgmt.FlagAlertingProvenanceLockWrites) {
|
|
return st.setProvenanceWithLocking(sess, recordKey, recordType, org, p)
|
|
}
|
|
return st.setProvenanceUpsert(sess, recordKey, recordType, org, p)
|
|
})
|
|
}
|
|
|
|
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
|
|
func (st DBstore) DeleteProvenance(ctx context.Context, o models.Provisionable, org int64) error {
|
|
return st.SQLStore.WithTransactionalDbSession(ctx, func(sess *db.Session) error {
|
|
_, err := sess.Delete(provenanceRecord{
|
|
RecordKey: o.ResourceID(),
|
|
RecordType: o.ResourceType(),
|
|
OrgID: org,
|
|
})
|
|
return err
|
|
})
|
|
}
|