diff --git a/pkg/registry/apps/alerting/notifications/receiver/authorize.go b/pkg/registry/apps/alerting/notifications/receiver/authorize.go index 60ac9c01cc4..3998201e748 100644 --- a/pkg/registry/apps/alerting/notifications/receiver/authorize.go +++ b/pkg/registry/apps/alerting/notifications/receiver/authorize.go @@ -10,6 +10,7 @@ import ( "github.com/grafana/grafana/pkg/apimachinery/errutil" "github.com/grafana/grafana/pkg/apimachinery/identity" "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/receiver/v0alpha1" + "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/receiver/v0alpha2" "github.com/grafana/grafana/pkg/services/ngalert/accesscontrol" ) @@ -23,7 +24,7 @@ type AccessControlService interface { } func Authorize(ctx context.Context, ac AccessControlService, attr authorizer.Attributes) (authorized authorizer.Decision, reason string, err error) { - if attr.GetResource() != v0alpha1.ResourceInfo.GroupResource().Resource { + if attr.GetResource() != v0alpha1.ResourceInfo.GroupResource().Resource && attr.GetResource() != v0alpha2.ResourceInfo.GroupResource().Resource { return authorizer.DecisionNoOpinion, "", nil } user, err := identity.GetRequester(ctx) diff --git a/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions.go b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions.go new file mode 100644 index 00000000000..3e8925e2e39 --- /dev/null +++ b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions.go @@ -0,0 +1,611 @@ +package v0alpha2 + +import ( + "errors" + "fmt" + "strings" + "unsafe" + + "github.com/grafana/alerting/receivers" + + jsoniter "github.com/json-iterator/go" + "github.com/modern-go/reflect2" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/types" + + model "github.com/grafana/grafana/apps/alerting/notifications/pkg/apis/receiver/v0alpha2" + "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" + gapiutil "github.com/grafana/grafana/pkg/services/apiserver/utils" + ngmodels "github.com/grafana/grafana/pkg/services/ngalert/models" + "github.com/grafana/grafana/pkg/util" +) + +func convertToK8sResources( + orgID int64, + receivers []*ngmodels.Receiver, + accesses map[string]ngmodels.ReceiverPermissionSet, + metadatas map[string]ngmodels.ReceiverMetadata, + namespacer request.NamespaceMapper, + selector fields.Selector, +) (*model.ReceiverList, error) { + result := &model.ReceiverList{ + Items: make([]model.Receiver, 0, len(receivers)), + } + for _, receiver := range receivers { + var access *ngmodels.ReceiverPermissionSet + if accesses != nil { + if a, ok := accesses[receiver.GetUID()]; ok { + access = &a + } + } + var metadata *ngmodels.ReceiverMetadata + if metadatas != nil { + if m, ok := metadatas[receiver.GetUID()]; ok { + metadata = &m + } + } + k8sResource, err := convertToK8sResource(orgID, receiver, access, metadata, namespacer) + if err != nil { + return nil, err + } + if selector != nil && !selector.Empty() && !selector.Matches(model.SelectableFields(k8sResource)) { + continue + } + result.Items = append(result.Items, *k8sResource) + } + return result, nil +} + +func convertToK8sResource( + orgID int64, + receiver *ngmodels.Receiver, + access *ngmodels.ReceiverPermissionSet, + metadata *ngmodels.ReceiverMetadata, + namespacer request.NamespaceMapper, +) (*model.Receiver, error) { + spec, err := specFromDomainReceiver(receiver) + if err != nil { + return nil, err + } + r := &model.Receiver{ + ObjectMeta: metav1.ObjectMeta{ + UID: types.UID(receiver.GetUID()), // This is needed to make PATCH work + Name: receiver.GetUID(), + Namespace: namespacer(orgID), + ResourceVersion: receiver.Version, + }, + Spec: spec, + } + r.SetProvenanceStatus(string(receiver.Provenance)) + + if access != nil { + for _, action := range ngmodels.ReceiverPermissions() { + mappedAction, ok := permissionMapper[action] + if !ok { + return nil, fmt.Errorf("unknown action %v", action) + } + if can, _ := access.Has(action); can { + r.SetAccessControl(mappedAction) + } + } + } + + if metadata != nil { + rules := make([]string, 0, len(metadata.InUseByRules)) + for _, rule := range metadata.InUseByRules { + rules = append(rules, rule.UID) + } + r.SetInUse(metadata.InUseByRoutes, rules) + } + r.UID = gapiutil.CalculateClusterWideUID(r) + return r, nil +} + +var permissionMapper = map[ngmodels.ReceiverPermission]string{ + ngmodels.ReceiverPermissionReadSecret: "canReadSecrets", + ngmodels.ReceiverPermissionAdmin: "canAdmin", + ngmodels.ReceiverPermissionWrite: "canWrite", + ngmodels.ReceiverPermissionDelete: "canDelete", +} + +// ContactPointFromContactPointExport parses the database model of the contact point (group of integrations) where settings are represented in JSON, +// to strongly typed ContactPoint. +func specFromDomainReceiver(domain *ngmodels.Receiver) (model.Spec, error) { + j := jsoniter.ConfigCompatibleWithStandardLibrary + j.RegisterExtension(&contactPointsExtension{}) + + result := model.Spec{ + Title: domain.Name, + } + + var errs []error + for _, rawIntegration := range domain.Integrations { + err := parseIntegration(j, &result, rawIntegration) + if err != nil { + // accumulate errors to report all at once. + errs = append(errs, fmt.Errorf("failed to parse %s integration (uid:%s): %w", rawIntegration.Config.Type, rawIntegration.UID, err)) + } + } + return result, errors.Join(errs...) +} + +// ContactPointToContactPointExport converts v0alpha2.Receiver to models.Receiver. +// It uses a special extension for json-iterator API that properly handles marshaling of some specific fields. +// +//nolint:gocyclo +func convertToDomainModel(apiModel *model.Receiver) (*ngmodels.Receiver, error) { + j := jsoniter.ConfigCompatibleWithStandardLibrary + // use json iterator with custom extension that has special codec for some field. + // This is needed to keep the API models clean and convert from a database model + j.RegisterExtension(&contactPointsExtension{}) + + cp := apiModel.Spec + + integration := make([]*ngmodels.Integration, 0, apiModel.Spec.IntegrationsCount()) + + var errs []error + for _, i := range cp.Alertmanager { + el, err := marshallIntegration(j, "prometheus-alertmanager", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Dingding { + el, err := marshallIntegration(j, "dingding", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Discord { + el, err := marshallIntegration(j, "discord", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Email { + el, err := marshallIntegration(j, "email", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Googlechat { + el, err := marshallIntegration(j, "googlechat", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Jira { + el, err := marshallIntegration(j, "jira", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Kafka { + el, err := marshallIntegration(j, "kafka", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Line { + el, err := marshallIntegration(j, "line", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Mqtt { + el, err := marshallIntegration(j, "mqtt", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Opsgenie { + el, err := marshallIntegration(j, "opsgenie", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Pagerduty { + el, err := marshallIntegration(j, "pagerduty", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Oncall { + el, err := marshallIntegration(j, "oncall", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Pushover { + el, err := marshallIntegration(j, "pushover", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Sensugo { + el, err := marshallIntegration(j, "sensugo", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Sns { + el, err := marshallIntegration(j, "sns", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Slack { + el, err := marshallIntegration(j, "slack", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Teams { + el, err := marshallIntegration(j, "teams", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Telegram { + el, err := marshallIntegration(j, "telegram", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Threema { + el, err := marshallIntegration(j, "threema", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Victorops { + el, err := marshallIntegration(j, "victorops", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Webhook { + el, err := marshallIntegration(j, "webhook", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Wecom { + el, err := marshallIntegration(j, "wecom", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + for _, i := range cp.Webex { + el, err := marshallIntegration(j, "webex", i, i.DisableResolveMessage, i.Uid, cp.Title) + if err != nil { + errs = append(errs, err) + } + integration = append(integration, el) + } + + if len(errs) > 0 { + return nil, errors.Join(errs...) + } + result := &ngmodels.Receiver{ + UID: apiModel.Name, + Name: apiModel.Spec.Title, + Integrations: integration, + Provenance: ngmodels.Provenance(apiModel.GetProvenanceStatus()), + Version: apiModel.ResourceVersion, + } + return result, nil +} + +// marshallIntegration converts the API model integration to the storage model that contains settings in the JSON format. +// The secret fields are not encrypted. +func marshallIntegration(json jsoniter.API, integrationType string, integration interface{}, disableResolveMessage *bool, uid *string, name string) (*ngmodels.Integration, error) { + data, err := json.Marshal(integration) + if err != nil { + return nil, fmt.Errorf("failed to marshall integration '%s' to JSON: %w", integrationType, err) + } + settings := map[string]interface{}{} + err = json.Unmarshal(data, &settings) + delete(settings, "uid") // integration UID is part of the integration + if err != nil { + return nil, fmt.Errorf("failed to marshall integration '%s' to map: %w", integrationType, err) + } + + config, err := ngmodels.IntegrationConfigFromType(integrationType) + if err != nil { + return nil, err + } + + e := &ngmodels.Integration{ + UID: "", + Name: name, + Config: config, + Settings: settings, + SecureSettings: nil, + } + if uid != nil { + e.UID = *uid + } + if disableResolveMessage != nil { + e.DisableResolveMessage = *disableResolveMessage + } + return e, nil +} + +//nolint:gocyclo +func parseIntegration(json jsoniter.API, result *model.Spec, integration *ngmodels.Integration) error { + var err error + var disable *bool + if integration.DisableResolveMessage { // populate only if true + disable = util.Pointer(integration.DisableResolveMessage) + } + + data, err := json.Marshal(integration.Settings) + if err != nil { + return fmt.Errorf("failed to marshall integration '%s' to JSON: %w", integration.Config.Type, err) + } + switch strings.ToLower(integration.Config.Type) { + case "prometheus-alertmanager": + integration := model.AlertmanagerIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Alertmanager = append(result.Alertmanager, integration) + } + case "dingding": + integration := model.DingdingIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Dingding = append(result.Dingding, integration) + } + case "discord": + integration := model.DiscordIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Discord = append(result.Discord, integration) + } + case "email": + integration := model.EmailIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Email = append(result.Email, integration) + } + case "googlechat": + integration := model.GooglechatIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Googlechat = append(result.Googlechat, integration) + } + case "jira": + integration := model.JiraIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Jira = append(result.Jira, integration) + } + case "kafka": + integration := model.KafkaIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Kafka = append(result.Kafka, integration) + } + case "line": + integration := model.LineIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Line = append(result.Line, integration) + } + case "mqtt": + integration := model.MqttIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Mqtt = append(result.Mqtt, integration) + } + case "opsgenie": + integration := model.OpsgenieIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Opsgenie = append(result.Opsgenie, integration) + } + case "pagerduty": + integration := model.PagerdutyIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Pagerduty = append(result.Pagerduty, integration) + } + case "oncall": + integration := model.OnCallIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Oncall = append(result.Oncall, integration) + } + case "pushover": + integration := model.PushoverIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Pushover = append(result.Pushover, integration) + } + case "sensugo": + integration := model.SensugoIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Sensugo = append(result.Sensugo, integration) + } + case "sns": + integration := model.SnsIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Sns = append(result.Sns, integration) + } + case "slack": + integration := model.SlackIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Slack = append(result.Slack, integration) + } + case "teams": + integration := model.TeamsIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Teams = append(result.Teams, integration) + } + case "telegram": + integration := model.TelegramIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Telegram = append(result.Telegram, integration) + } + case "threema": + integration := model.ThreemaIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Threema = append(result.Threema, integration) + } + case "victorops": + integration := model.VictoropsIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Victorops = append(result.Victorops, integration) + } + case "webhook": + integration := model.WebhookIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Webhook = append(result.Webhook, integration) + } + case "wecom": + integration := model.WecomIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Wecom = append(result.Wecom, integration) + } + case "webex": + integration := model.WebexIntegration{DisableResolveMessage: disable, Uid: util.Pointer(integration.UID)} + if err = json.Unmarshal(data, &integration); err == nil { + result.Webex = append(result.Webex, integration) + } + default: + err = fmt.Errorf("integration %s is not supported", integration.Config.Type) + } + return err +} + +// contactPointsExtension extends jsoniter with special codecs for some integrations' fields that are encoded differently in the legacy configuration. +type contactPointsExtension struct { + jsoniter.DummyExtension +} + +func (c contactPointsExtension) UpdateStructDescriptor(structDescriptor *jsoniter.StructDescriptor) { + if structDescriptor.Type == reflect2.TypeOf(model.EmailIntegration{}) { + bind := structDescriptor.GetField("Addresses") + codec := &emailAddressCodec{} + bind.Decoder = codec + bind.Encoder = codec + } + if structDescriptor.Type == reflect2.TypeOf(model.PushoverIntegration{}) { + codec := &numberAsStringCodec{} + for _, field := range []string{"Priority", "OkPriority"} { + desc := structDescriptor.GetField(field) + desc.Decoder = codec + desc.Encoder = codec + } + // the same logic is in the pushover.NewConfig in alerting module + codec = &numberAsStringCodec{ignoreError: true} + for _, field := range []string{"Retry", "Expire"} { + desc := structDescriptor.GetField(field) + desc.Decoder = codec + desc.Encoder = codec + } + } + if structDescriptor.Type == reflect2.TypeOf(model.WebhookIntegration{}) { + codec := &numberAsStringCodec{ignoreError: true} + desc := structDescriptor.GetField("MaxAlerts") + desc.Decoder = codec + desc.Encoder = codec + } + if structDescriptor.Type == reflect2.TypeOf(model.OnCallIntegration{}) { + codec := &numberAsStringCodec{ignoreError: true} + desc := structDescriptor.GetField("MaxAlerts") + desc.Decoder = codec + desc.Encoder = codec + } + if structDescriptor.Type == reflect2.TypeOf(model.MqttIntegration{}) { + codec := &numberAsStringCodec{ignoreError: true} + desc := structDescriptor.GetField("Qos") + desc.Decoder = codec + desc.Encoder = codec + } +} + +type emailAddressCodec struct{} + +func (d *emailAddressCodec) IsEmpty(ptr unsafe.Pointer) bool { + f := *(*[]string)(ptr) + return len(f) == 0 +} + +func (d *emailAddressCodec) Encode(ptr unsafe.Pointer, stream *jsoniter.Stream) { + f := *(*[]string)(ptr) + addresses := strings.Join(f, ";") + stream.WriteString(addresses) +} + +func (d *emailAddressCodec) Decode(ptr unsafe.Pointer, iter *jsoniter.Iterator) { + s := iter.ReadString() + emails := strings.FieldsFunc(strings.Trim(s, "\""), func(r rune) bool { + switch r { + case ',', ';', '\n': + return true + } + return false + }) + *((*[]string)(ptr)) = emails +} + +// converts a string representation of a number to *int64 +type numberAsStringCodec struct { + ignoreError bool // if true, then ignores the error and keeps value nil +} + +func (d *numberAsStringCodec) IsEmpty(ptr unsafe.Pointer) bool { + return *((*(*int))(ptr)) == nil +} + +func (d *numberAsStringCodec) Encode(ptr unsafe.Pointer, stream *jsoniter.Stream) { + val := *((*(*int))(ptr)) + if val == nil { + stream.WriteNil() + return + } + stream.WriteInt(*val) +} + +func (d *numberAsStringCodec) Decode(ptr unsafe.Pointer, iter *jsoniter.Iterator) { + valueType := iter.WhatIsNext() + var value int64 + switch valueType { + case jsoniter.NumberValue: + value = iter.ReadInt64() + case jsoniter.StringValue: + var num receivers.OptionalNumber + err := num.UnmarshalJSON(iter.ReadStringAsSlice()) + if err != nil { + iter.ReportError("numberAsStringCodec", fmt.Sprintf("failed to unmarshall string as OptionalNumber: %s", err.Error())) + } + if num.String() == "" { + return + } + value, err = num.Int64() + if err != nil { + if !d.ignoreError { + iter.ReportError("numberAsStringCodec", fmt.Sprintf("string does not represent an integer number: %s", err.Error())) + } + return + } + case jsoniter.NilValue: + iter.ReadNil() + return + default: + iter.ReportError("numberAsStringCodec", "not number or string") + } + *((*(*int64))(ptr)) = &value +} diff --git a/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions_test.go b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions_test.go new file mode 100644 index 00000000000..2ffe9decd17 --- /dev/null +++ b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/conversions_test.go @@ -0,0 +1,60 @@ +package v0alpha2 + +import ( + "encoding/json" + "fmt" + "testing" + + "github.com/google/go-cmp/cmp" + "github.com/grafana/alerting/notify" + "github.com/stretchr/testify/require" + + "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" + "github.com/grafana/grafana/pkg/services/ngalert/models" +) + +func TestConvertToK8sResource(t *testing.T) { + for integrationType, cfg := range notify.AllKnownConfigsForTesting { + t.Run(integrationType, func(t *testing.T) { + schema, err := models.IntegrationConfigFromType(integrationType) + require.NoError(t, err) + settingsMap := map[string]any{} + require.NoError(t, json.Unmarshal([]byte(cfg.Config), &settingsMap)) + recCfg := &models.Receiver{ + UID: fmt.Sprintf("uid-%s", integrationType), + Name: fmt.Sprintf("name-%s", integrationType), + Integrations: []*models.Integration{ + { + UID: fmt.Sprintf("intuid-%s", integrationType), + Name: fmt.Sprintf("name-%s", integrationType), + Config: schema, + DisableResolveMessage: false, + Settings: settingsMap, + SecureSettings: nil, + }, + }, + Provenance: "api", + Version: "1234", + } + + result, err := convertToK8sResource(1, recCfg, &models.ReceiverPermissionSet{}, &models.ReceiverMetadata{}, request.GetNamespaceMapper(nil)) + require.NoError(t, err) + + back, err := convertToDomainModel(result) + require.NoError(t, err) + + switch integrationType { + case "webhook": // webhook has a different format for maxAlerts, in test config it's a string, in API model it's an int. + val := back.Integrations[0].Settings["maxAlerts"].(float64) + back.Integrations[0].Settings["maxAlerts"] = fmt.Sprintf("%d", int(val)) + case "mqtt": // same for qos + val := back.Integrations[0].Settings["qos"].(float64) + back.Integrations[0].Settings["qos"] = fmt.Sprintf("%d", int(val)) + } + diff := cmp.Diff(recCfg, back) + if len(diff) != 0 { + require.Failf(t, "The re-marshalled configuration does not match the expected one", diff) + } + }) + } +} diff --git a/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/legacy_storage.go b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/legacy_storage.go new file mode 100644 index 00000000000..8888525da8b --- /dev/null +++ b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/legacy_storage.go @@ -0,0 +1,279 @@ +package v0alpha2 + +import ( + "context" + "errors" + "fmt" + + apierrors "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/apis/meta/internalversion" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apiserver/pkg/registry/rest" + + model "github.com/grafana/grafana/apps/alerting/notifications/pkg/apis/receiver/v0alpha2" + "github.com/grafana/grafana/pkg/apimachinery/identity" + grafanarest "github.com/grafana/grafana/pkg/apiserver/rest" + "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" + alertingac "github.com/grafana/grafana/pkg/services/ngalert/accesscontrol" + "github.com/grafana/grafana/pkg/services/ngalert/api/tooling/definitions" + ngmodels "github.com/grafana/grafana/pkg/services/ngalert/models" + "github.com/grafana/grafana/pkg/services/ngalert/notifier/legacy_storage" +) + +var ( + _ grafanarest.Storage = (*legacyStorage)(nil) +) + +type ReceiverService interface { + GetReceiver(ctx context.Context, q ngmodels.GetReceiverQuery, user identity.Requester) (*ngmodels.Receiver, error) + GetReceivers(ctx context.Context, q ngmodels.GetReceiversQuery, user identity.Requester) ([]*ngmodels.Receiver, error) + CreateReceiver(ctx context.Context, r *ngmodels.Receiver, orgID int64, user identity.Requester) (*ngmodels.Receiver, error) + UpdateReceiver(ctx context.Context, r *ngmodels.Receiver, storedSecureFields map[string][]string, orgID int64, user identity.Requester) (*ngmodels.Receiver, error) + DeleteReceiver(ctx context.Context, name string, provenance definitions.Provenance, version string, orgID int64, user identity.Requester) error +} + +type MetadataService interface { + AccessControlMetadata(ctx context.Context, user identity.Requester, receivers ...*ngmodels.Receiver) (map[string]ngmodels.ReceiverPermissionSet, error) + InUseMetadata(ctx context.Context, orgID int64, receivers ...*ngmodels.Receiver) (map[string]ngmodels.ReceiverMetadata, error) +} + +type legacyStorage struct { + service ReceiverService + namespacer request.NamespaceMapper + tableConverter rest.TableConvertor + metadata MetadataService +} + +func (s *legacyStorage) New() runtime.Object { + return ResourceInfo.NewFunc() +} + +func (s *legacyStorage) Destroy() {} + +func (s *legacyStorage) NamespaceScoped() bool { + return true // namespace == org +} + +func (s *legacyStorage) GetSingularName() string { + return ResourceInfo.GetSingularName() +} + +func (s *legacyStorage) NewList() runtime.Object { + return ResourceInfo.NewListFunc() +} + +func (s *legacyStorage) ConvertToTable(ctx context.Context, object runtime.Object, tableOptions runtime.Object) (*metav1.Table, error) { + return s.tableConverter.ConvertToTable(ctx, object, tableOptions) +} + +func (s *legacyStorage) List(ctx context.Context, opts *internalversion.ListOptions) (runtime.Object, error) { + orgId, err := request.OrgIDForList(ctx) + if err != nil { + return nil, err + } + + q := ngmodels.GetReceiversQuery{ + OrgID: orgId, + Decrypt: true, + } + + user, err := identity.GetRequester(ctx) + if err != nil { + return nil, err + } + + res, err := s.service.GetReceivers(ctx, q, user) + if err != nil { + // This API should not be returning a forbidden error when the user does not have access to any resources. + // This can be true for a contact point creator role, for example. + // This should eventually be changed downstream in the auth logic but provisioning API currently relies on this + // behaviour to return useful forbidden errors when exporting decrypted receivers. + if !errors.Is(err, alertingac.ErrAuthorizationBase) { + return nil, err + } + res = nil + } + + accesses, err := s.metadata.AccessControlMetadata(ctx, user, res...) + if err != nil { + return nil, fmt.Errorf("failed to get access control metadata: %w", err) + } + + inUses, err := s.metadata.InUseMetadata(ctx, orgId, res...) + if err != nil { + return nil, fmt.Errorf("failed to get in-use metadata: %w", err) + } + + return convertToK8sResources(orgId, res, accesses, inUses, s.namespacer, opts.FieldSelector) +} + +func (s *legacyStorage) Get(ctx context.Context, uid string, _ *metav1.GetOptions) (runtime.Object, error) { + info, err := request.NamespaceInfoFrom(ctx, true) + if err != nil { + return nil, err + } + + name, err := legacy_storage.UidToName(uid) + if err != nil { + return nil, apierrors.NewNotFound(ResourceInfo.GroupResource(), uid) + } + q := ngmodels.GetReceiverQuery{ + OrgID: info.OrgID, + Name: name, + Decrypt: false, + } + + user, err := identity.GetRequester(ctx) + if err != nil { + return nil, err + } + + r, err := s.service.GetReceiver(ctx, q, user) + if err != nil { + return nil, err + } + + var access *ngmodels.ReceiverPermissionSet + accesses, err := s.metadata.AccessControlMetadata(ctx, user, r) + if err == nil { + if a, ok := accesses[r.GetUID()]; ok { + access = &a + } + } else { + return nil, fmt.Errorf("failed to get access control metadata: %w", err) + } + + var inUse *ngmodels.ReceiverMetadata + inUses, err := s.metadata.InUseMetadata(ctx, info.OrgID, r) + if err == nil { + if a, ok := inUses[r.GetUID()]; ok { + inUse = &a + } + } else { + return nil, fmt.Errorf("failed to get access control metadata: %w", err) + } + + return convertToK8sResource(info.OrgID, r, access, inUse, s.namespacer) +} + +func (s *legacyStorage) Create(ctx context.Context, + obj runtime.Object, + createValidation rest.ValidateObjectFunc, + _ *metav1.CreateOptions, +) (runtime.Object, error) { + info, err := request.NamespaceInfoFrom(ctx, true) + if err != nil { + return nil, err + } + if createValidation != nil { + if err := createValidation(ctx, obj.DeepCopyObject()); err != nil { + return nil, err + } + } + p, ok := obj.(*model.Receiver) + if !ok { + return nil, fmt.Errorf("expected receiver but got %s", obj.GetObjectKind().GroupVersionKind()) + } + if p.Name != "" { // TODO remove when metadata.name can be defined by user + return nil, apierrors.NewBadRequest("object's metadata.name should be empty") + } + model, err := convertToDomainModel(p) + if err != nil { + return nil, err + } + + user, err := identity.GetRequester(ctx) + if err != nil { + return nil, err + } + + out, err := s.service.CreateReceiver(ctx, model, info.OrgID, user) + if err != nil { + return nil, err + } + return convertToK8sResource(info.OrgID, out, nil, nil, s.namespacer) +} + +func (s *legacyStorage) Update(ctx context.Context, + uid string, + objInfo rest.UpdatedObjectInfo, + createValidation rest.ValidateObjectFunc, + updateValidation rest.ValidateObjectUpdateFunc, + _ bool, + _ *metav1.UpdateOptions, +) (runtime.Object, bool, error) { + info, err := request.NamespaceInfoFrom(ctx, true) + if err != nil { + return nil, false, err + } + + user, err := identity.GetRequester(ctx) + if err != nil { + return nil, false, err + } + + old, err := s.Get(ctx, uid, nil) + if err != nil { + return old, false, err + } + obj, err := objInfo.UpdatedObject(ctx, old) + if err != nil { + return old, false, err + } + if updateValidation != nil { + if err := updateValidation(ctx, obj, old); err != nil { + return nil, false, err + } + } + p, ok := obj.(*model.Receiver) + if !ok { + return nil, false, fmt.Errorf("expected receiver but got %s", obj.GetObjectKind().GroupVersionKind()) + } + model, err := convertToDomainModel(p) + if err != nil { + return old, false, err + } + + updated, err := s.service.UpdateReceiver(ctx, model, nil, info.OrgID, user) + if err != nil { + return nil, false, err + } + + r, err := convertToK8sResource(info.OrgID, updated, nil, nil, s.namespacer) + return r, false, err +} + +// GracefulDeleter +func (s *legacyStorage) Delete(ctx context.Context, uid string, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions) (runtime.Object, bool, error) { + info, err := request.NamespaceInfoFrom(ctx, true) + if err != nil { + return nil, false, err + } + + user, err := identity.GetRequester(ctx) + if err != nil { + return nil, false, err + } + + old, err := s.Get(ctx, uid, nil) + if err != nil { + return old, false, err + } + if deleteValidation != nil { + if err = deleteValidation(ctx, old); err != nil { + return nil, false, err + } + } + version := "" + if options.Preconditions != nil && options.Preconditions.ResourceVersion != nil { + version = *options.Preconditions.ResourceVersion + } + + err = s.service.DeleteReceiver(ctx, uid, definitions.Provenance(ngmodels.ProvenanceNone), version, info.OrgID, user) // TODO add support for dry-run option + return old, false, err // false - will be deleted async +} + +func (s *legacyStorage) DeleteCollection(ctx context.Context, deleteValidation rest.ValidateObjectFunc, options *metav1.DeleteOptions, listOptions *internalversion.ListOptions) (runtime.Object, error) { + return nil, apierrors.NewMethodNotSupported(ResourceInfo.GroupResource(), "deleteCollection") +} diff --git a/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/storage.go b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/storage.go new file mode 100644 index 00000000000..2c1b339c908 --- /dev/null +++ b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/storage.go @@ -0,0 +1,19 @@ +package v0alpha2 + +import ( + grafanarest "github.com/grafana/grafana/pkg/apiserver/rest" + "github.com/grafana/grafana/pkg/services/apiserver/endpoints/request" +) + +func NewStorage( + legacySvc ReceiverService, + namespacer request.NamespaceMapper, + metadata MetadataService, +) grafanarest.Storage { + return &legacyStorage{ + service: legacySvc, + namespacer: namespacer, + tableConverter: ResourceInfo.TableConverter(), + metadata: metadata, + } +} diff --git a/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/type.go b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/type.go new file mode 100644 index 00000000000..e28095fd539 --- /dev/null +++ b/pkg/registry/apps/alerting/notifications/receiver/v0alpha2/type.go @@ -0,0 +1,38 @@ +package v0alpha2 + +import ( + "fmt" + "strings" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + + modelv2 "github.com/grafana/grafana/apps/alerting/notifications/pkg/apis/receiver/v0alpha2" + "github.com/grafana/grafana/pkg/apimachinery/utils" +) + +var kindV2 = modelv2.Kind() + +var ResourceInfo = utils.NewResourceInfo(kindV2.Group(), kindV2.Version(), + kindV2.GroupVersionResource().Resource, strings.ToLower(kindV2.Kind()), kindV2.Kind(), + func() runtime.Object { return kindV2.ZeroValue() }, + func() runtime.Object { return kindV2.ZeroListValue() }, + utils.TableColumns{ + Definition: []metav1.TableColumnDefinition{ + {Name: "Name", Type: "string", Format: "name"}, + {Name: "Title", Type: "string", Format: "string", Description: "The receiver name"}, + {Name: "Integrations", Type: "string", Format: "string", Description: "The integration types"}, + }, + Reader: func(obj any) ([]interface{}, error) { + r, ok := obj.(*modelv2.Receiver) + if ok { + return []interface{}{ + r.Name, + r.Spec.Title, + strings.Join(r.Spec.GetIntegrationsTypes(), ","), + }, nil + } + return nil, fmt.Errorf("expected resource or info") + }, + }, +) diff --git a/pkg/registry/apps/alerting/notifications/register.go b/pkg/registry/apps/alerting/notifications/register.go index 33d5242b462..cf3bb340a13 100644 --- a/pkg/registry/apps/alerting/notifications/register.go +++ b/pkg/registry/apps/alerting/notifications/register.go @@ -13,6 +13,7 @@ import ( grafanarest "github.com/grafana/grafana/pkg/apiserver/rest" "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/receiver" receiverv1 "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/receiver/v0alpha1" + receiverv2 "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/receiver/v0alpha2" "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/routingtree" "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/templategroup" "github.com/grafana/grafana/pkg/registry/apps/alerting/notifications/timeinterval" @@ -56,7 +57,8 @@ func getAuthorizer(authz accesscontrol.AccessControl) authorizer.Authorizer { return templategroup.Authorize(ctx, authz, a) case timeinterval.ResourceInfo.GroupResource().Resource: return timeinterval.Authorize(ctx, authz, a) - case receiverv1.ResourceInfo.GroupResource().Resource: + case receiverv1.ResourceInfo.GroupResource().Resource, + receiverv2.ResourceInfo.GroupResource().Resource: return receiver.Authorize(ctx, ac.NewReceiverAccess[*ngmodels.Receiver](authz, false), a) case routingtree.ResourceInfo.GroupResource().Resource: return routingtree.Authorize(ctx, authz, a) @@ -75,6 +77,8 @@ func getLegacyStorage(namespacer request.NamespaceMapper, ng *ngalert.AlertNG) r return templategroup.NewStorage(ng.Api.Templates, namespacer) } else if gvr == routingtree.ResourceInfo.GroupVersionResource() { return routingtree.NewStorage(ng.Api.Policies, namespacer) + } else if gvr == receiverv2.ResourceInfo.GroupVersionResource() { + return receiverv2.NewStorage(ng.Api.ReceiverService, namespacer, ng.Api.ReceiverService) } panic("unknown legacy storage requested: " + gvr.String()) }