From 37049bd5f3fcde4d4241e8ec45c598c340915141 Mon Sep 17 00:00:00 2001 From: Rain Date: Wed, 19 Feb 2020 19:47:18 +0800 Subject: [PATCH] feature(notify): Validate config --- cmd/climc/shell/notify.go | 19 +++- pkg/notify/compatible_config.go | 64 +----------- pkg/notify/dispatcher.go | 77 ++++++++++++--- pkg/notify/handlers.go | 18 ++++ pkg/notify/interface/interface.go | 1 + pkg/notify/models/mod_config.go | 81 ++++++++++++++++ pkg/notify/rpc/apis/send_server.pb.go | 134 +++++++++++++++++++++----- pkg/notify/rpc/apis/send_server.proto | 8 +- pkg/notify/rpc/send.go | 44 ++++++++- 9 files changed, 344 insertions(+), 102 deletions(-) diff --git a/cmd/climc/shell/notify.go b/cmd/climc/shell/notify.go index 8741bd701b..69d0d315c6 100644 --- a/cmd/climc/shell/notify.go +++ b/cmd/climc/shell/notify.go @@ -34,7 +34,24 @@ func init() { } body := jsonutils.NewDict() body.Add(tmp, args.CONTACTTYPE) - modules.Configs.Create(s, body) + ret, err := modules.Configs.Create(s, body) + if err != nil { + return err + } + printObject(ret) + return nil + }) + R(&ConfigCreate2Options{}, "notify-config-validate", "config validate", func(s *mcclient.ClientSession, + args *ConfigCreate2Options) error { + tmp := jsonutils.NewDict() + for i := 0; i+1 < len(args.CONFIGS); i += 2 { + tmp.Add(jsonutils.NewString(args.CONFIGS[i+1]), args.CONFIGS[i]) + } + ret, err := modules.Configs.PerformAction(s, args.CONTACTTYPE, "validate", tmp) + if err != nil { + return err + } + printObject(ret) return nil }) diff --git a/pkg/notify/compatible_config.go b/pkg/notify/compatible_config.go index 074ba5b209..1412fbf6be 100644 --- a/pkg/notify/compatible_config.go +++ b/pkg/notify/compatible_config.go @@ -17,7 +17,6 @@ package notify import ( "context" "net/http" - "strconv" "yunion.io/x/jsonutils" @@ -56,74 +55,19 @@ func emailConfigGetHandler(ctx context.Context, w http.ResponseWriter, r *http.R if err != nil { httperrors.GeneralServerError(w, err) } - // hostport should be int and ssl_global should be bool - newDataDict := make(map[string]interface{}) data, _ := ret.Get("config") dataDict := data.(*jsonutils.JSONDict) - dataDict = database2Display(dataDict) - for _, k := range dataDict.SortedKeys() { - tmp, _ := dataDict.GetString(k) - switch k { - case "hostport": - port, _ := strconv.Atoi(tmp) - newDataDict[k] = port - case "ssl_global": - ssl, _ := strconv.ParseBool(tmp) - newDataDict[k] = ssl - default: - newDataDict[k] = tmp - } - } - appsrv.SendJSON(w, jsonutils.Marshal(map[string]map[string]interface{}{ - EMAIL_KEYWORD: newDataDict, - })) -} - -func dispaly2Database(dict *jsonutils.JSONDict) *jsonutils.JSONDict { - keys := dict.SortedKeys() - newKey := "" - for _, key := range keys { - switch key { - case "username", "password": - newKey = "mail." + key - case "hostname", "hostport": - newKey = "mail.smtp." + key - case "ssl_global": - newKey = "mail.global.ssl" - } - v, _ := dict.Get(key) - dict.Add(v, newKey) - dict.Remove(key) - } - return dict -} - -func database2Display(dict *jsonutils.JSONDict) *jsonutils.JSONDict { - keys := dict.SortedKeys() - newKey := "" - for _, key := range keys { - switch key { - case "mail.username", "mail.password": - newKey = key[5:] - case "mail.smtp.hostname", "mail.smtp.hostport": - newKey = key[10:] - case "mail.global.ssl": - newKey = "ssl_global" - } - v, _ := dict.Get(key) - dict.Add(v, newKey) - dict.Remove(key) - } - return dict + output := jsonutils.NewDict() + output.Add(dataDict, EMAIL_KEYWORD) + appsrv.SendJSON(w, output) } func emailConfigUpdateHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { manager, _, _, body := fetchEnv(ctx, w, r) body, _ = body.Get(EMAIL_KEYWORD) bodyRet := jsonutils.DeepCopy(body) - bodyDict := body.(*jsonutils.JSONDict) newBody := jsonutils.NewDict() - newBody.Add(dispaly2Database(bodyDict), EMAIL) + newBody.Add(body, EMAIL) err := manager.UpdateConfig(ctx, newBody) if err != nil { httperrors.GeneralServerError(w, err) diff --git a/pkg/notify/dispatcher.go b/pkg/notify/dispatcher.go index 26791ef8f1..44d7186d7a 100644 --- a/pkg/notify/dispatcher.go +++ b/pkg/notify/dispatcher.go @@ -56,15 +56,19 @@ func (self *NotifyModelDispatcher) GetConfig(ctx context.Context, params map[str if err != nil { return nil, err } - keyVs := make(map[string]string) - for _, ret := range listResult.Data { - key, _ := ret.GetString("key_text") - value, _ := ret.GetString("value_text") - keyVs[key] = value + configs := jsonutils.NewDict() + for _, data := range listResult.Data { + key, _ := data.GetString("key_text") + value, _ := data.Get("value_text") + configs.Add(value, key) } - return jsonutils.Marshal(map[string]map[string]string{ - models.ConfigManager.Keyword(): keyVs, - }), nil + cType, ok := params[""] + if ok { + configs = models.ConfigManager.Database2Display(cType, configs) + } + output := jsonutils.NewDict() + output.Add(configs, models.ConfigManager.Keyword()) + return output, nil } func (self *NotifyModelDispatcher) DeleteConfig(ctx context.Context, params map[string]string) error { @@ -93,6 +97,7 @@ func (self *NotifyModelDispatcher) UpdateConfig(ctx context.Context, body jsonut } tmp, _ := data.Get(contactType) data = tmp.(*jsonutils.JSONDict) + data = models.ConfigManager.Display2Database(contactType, data) userCred := policy.FetchUserCredential(ctx) // If no config of type 'contactType' in database, create news. // Else delete original ones and create news. @@ -109,26 +114,74 @@ func (self *NotifyModelDispatcher) UpdateConfig(ctx context.Context, body jsonut } } } + keys := data.SortedKeys() config := make(map[string]string) - // create - for _, key := range data.SortedKeys() { + createDataList := make([]jsonutils.JSONObject, 0, len(keys)) + + // Extract data + for _, key := range keys { createData := jsonutils.NewDict() tmp, _ = data.Get(key) createData.Add(tmp, "value_text") createData.Add(jsonutils.NewString(key), "key_text") createData.Add(jsonutils.NewString(contactType), "type") - _, err := self.Create(ctx, jsonutils.JSONNull, createData, nil) + createDataList = append(createDataList, createData) value, _ := tmp.GetString() config[key] = value + } + + // validate configs + isValid, message, err := models.NotifyService.ValidateConfig(ctx, contactType, config) + if err != nil { + if errors.Cause(err) != errors.ErrNotImplemented { + return httperrors.NewInternalServerError("Validate Config error: %s", err.Error()) + } + isValid = true + } + if !isValid { + return httperrors.NewInputParameterError("validate failed: %s", message) + } + + // create + for _, createData := range createDataList { + _, err := self.Create(ctx, jsonutils.JSONNull, createData, nil) if err != nil { - return errors.Wrapf(err, "Create config (%s, %s, %s) failed", contactType, key, value) + return errors.Wrapf(err, "Create config %s for contact type %s failed", createData.String(), contactType) } } + // update config models.RestartService(config, contactType) return nil } +func (self *NotifyModelDispatcher) ValidateConfig(ctx context.Context, contactType string, body jsonutils.JSONObject) (jsonutils.JSONObject, error) { + dict, ok := body.(*jsonutils.JSONDict) + if !ok { + return nil, httperrors.NewInputParameterError("") + } + dict = models.ConfigManager.Display2Database(contactType, dict) + configs := make(map[string]string) + for _, key := range dict.SortedKeys() { + value, err := dict.GetString(key) + if err != nil { + return nil, errors.Wrap(err, "jsonutils.JsonDict.GetString") + } + configs[key] = value + } + isValid, message, err := models.NotifyService.ValidateConfig(ctx, contactType, configs) + if err != nil { + if errors.Cause(err) == errors.ErrNotImplemented { + return nil, httperrors.NewNotImplementedError("validating config of %s", contactType) + } + return nil, err + } + ret := jsonutils.NewDict() + ret.Add(jsonutils.NewBool(isValid), "is_valid") + ret.Add(jsonutils.NewString(message), "message") + return ret, nil +} + // CreateNotification create new notifications and send them through rpc.RpcService. // If data contains 'gid' field, that means that send message to all users in group. // Else send messager to user whose uid equals 'uid' in data. diff --git a/pkg/notify/handlers.go b/pkg/notify/handlers.go index 4949744aef..95b4f4d9cc 100644 --- a/pkg/notify/handlers.go +++ b/pkg/notify/handlers.go @@ -139,6 +139,9 @@ func AddNotifyDispatcher(prefix string, app *appsrv.Application) { app.AddHandler2("DELETE", fmt.Sprintf("%s/%s/", prefix, modelDispatcher.KeywordPlural()), middleware(configDeleteHandler), metadata, "delete_configs", tags) + app.AddHandler2("POST", + fmt.Sprintf("%s/%s//validate", prefix, modelDispatcher.KeywordPlural()), + middleware(configValidateHandler), metadata, "validate_configs", tags) // email handler for being compatible app.AddHandler2("POST", @@ -246,6 +249,21 @@ func configUpdateHandler(ctx context.Context, w http.ResponseWriter, r *http.Req } } +func configValidateHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { + manager, params, _, body := fetchEnv(ctx, w, r) + body, err := body.Get(models.ConfigManager.Keyword()) + if err != nil { + httperrors.GeneralServerError(w, httperrors.NewInputParameterError("need config")) + } + ctype := params[""] + res, err := manager.ValidateConfig(ctx, ctype, body) + if err != nil { + httperrors.GeneralServerError(w, err) + return + } + appsrv.SendJSON(w, res) +} + func notificationHandler(ctx context.Context, w http.ResponseWriter, r *http.Request) { manager, _, _, body := fetchEnv(ctx, w, r) data, err := body.Get(manager.Keyword()) diff --git a/pkg/notify/interface/interface.go b/pkg/notify/interface/interface.go index 44eaaccac5..21d72105e2 100644 --- a/pkg/notify/interface/interface.go +++ b/pkg/notify/interface/interface.go @@ -28,6 +28,7 @@ type INotifyService interface { RestartService(ctx context.Context, config SConfig, serviceName string) Send(ctx context.Context, contactType, contact, topic, msg, priority string) error ContactByMobile(ctx context.Context, mobile, serviceName string) (string, error) + ValidateConfig(ctx context.Context, cType string, configs map[string]string) (isValid bool, message string, err error) } type IServiceConfigStore interface { diff --git a/pkg/notify/models/mod_config.go b/pkg/notify/models/mod_config.go index e904df8ced..f135fca732 100644 --- a/pkg/notify/models/mod_config.go +++ b/pkg/notify/models/mod_config.go @@ -17,6 +17,7 @@ package models import ( "context" "fmt" + "strconv" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -169,3 +170,83 @@ func (self *SConfigManager) GetConfig(contactType string) (_interface.SConfig, e func (self *SConfigManager) SetConfig(contactType string, config _interface.SConfig) error { return fmt.Errorf("SetConfig Not Implemented") } + +type SConvertFunc func(*jsonutils.JSONDict) *jsonutils.JSONDict + +var ( + // toDisplay store SConvertFunc which convert the data from client to the form required by database of contactType + toDisplay map[string]SConvertFunc + // fromDisplay store SConvertFunc which convert the data from database to the form required by client of contactType + fromDisplay map[string]SConvertFunc +) + +func init() { + toDisplay = map[string]SConvertFunc{ + EMAIL: emailToDisplay, + } + fromDisplay = map[string]SConvertFunc{ + EMAIL: emailFromDisplay, + } +} + +func emailFromDisplay(dict *jsonutils.JSONDict) *jsonutils.JSONDict { + keys := dict.SortedKeys() + newKey := "" + for _, key := range keys { + switch key { + case "username", "password": + newKey = "mail." + key + case "hostname", "hostport": + newKey = "mail.smtp." + key + case "ssl_global": + newKey = "mail.global.ssl" + } + v, _ := dict.Get(key) + dict.Add(v, newKey) + dict.Remove(key) + } + return dict +} + +func emailToDisplay(dict *jsonutils.JSONDict) *jsonutils.JSONDict { + keys := dict.SortedKeys() + for _, key := range keys { + newKey := "" + value, _ := dict.Get(key) + switch key { + case "mail.username", "mail.password": + newKey = key[5:] + case "mail.smtp.hostname": + newKey = key[10:] + case "mail.smtp.hostport": + newKey = key[10:] + portStr, _ := value.GetString() + port, _ := strconv.Atoi(portStr) + value = jsonutils.NewInt(int64(port)) + case "mail.global.ssl": + newKey = "ssl_global" + sslStr, _ := value.GetString() + ssl, _ := strconv.ParseBool(sslStr) + value = jsonutils.NewBool(ssl) + } + dict.Add(value, newKey) + dict.Remove(key) + } + return dict +} + +func (self *SConfigManager) Display2Database(cType string, dict *jsonutils.JSONDict) *jsonutils.JSONDict { + cf, ok := fromDisplay[cType] + if ok { + return cf(dict) + } + return dict +} + +func (self *SConfigManager) Database2Display(cType string, dict *jsonutils.JSONDict) *jsonutils.JSONDict { + cf, ok := toDisplay[cType] + if ok { + return cf(dict) + } + return dict +} diff --git a/pkg/notify/rpc/apis/send_server.pb.go b/pkg/notify/rpc/apis/send_server.pb.go index f1db91f6b1..40ddf6bf3c 100644 --- a/pkg/notify/rpc/apis/send_server.pb.go +++ b/pkg/notify/rpc/apis/send_server.pb.go @@ -266,6 +266,53 @@ func (m *UseridByMobileReply) GetUserid() string { return "" } +type ValidateConfigReply struct { + IsValid bool `protobuf:"varint,1,opt,name=isValid,proto3" json:"isValid,omitempty"` + Msg string `protobuf:"bytes,2,opt,name=msg,proto3" json:"msg,omitempty"` + XXX_NoUnkeyedLiteral struct{} `json:"-"` + XXX_unrecognized []byte `json:"-"` + XXX_sizecache int32 `json:"-"` +} + +func (m *ValidateConfigReply) Reset() { *m = ValidateConfigReply{} } +func (m *ValidateConfigReply) String() string { return proto.CompactTextString(m) } +func (*ValidateConfigReply) ProtoMessage() {} +func (*ValidateConfigReply) Descriptor() ([]byte, []int) { + return fileDescriptor_63fdd68f7eb311f9, []int{5} +} + +func (m *ValidateConfigReply) XXX_Unmarshal(b []byte) error { + return xxx_messageInfo_ValidateConfigReply.Unmarshal(m, b) +} +func (m *ValidateConfigReply) XXX_Marshal(b []byte, deterministic bool) ([]byte, error) { + return xxx_messageInfo_ValidateConfigReply.Marshal(b, m, deterministic) +} +func (m *ValidateConfigReply) XXX_Merge(src proto.Message) { + xxx_messageInfo_ValidateConfigReply.Merge(m, src) +} +func (m *ValidateConfigReply) XXX_Size() int { + return xxx_messageInfo_ValidateConfigReply.Size(m) +} +func (m *ValidateConfigReply) XXX_DiscardUnknown() { + xxx_messageInfo_ValidateConfigReply.DiscardUnknown(m) +} + +var xxx_messageInfo_ValidateConfigReply proto.InternalMessageInfo + +func (m *ValidateConfigReply) GetIsValid() bool { + if m != nil { + return m.IsValid + } + return false +} + +func (m *ValidateConfigReply) GetMsg() string { + if m != nil { + return m.Msg + } + return "" +} + func init() { proto.RegisterType((*SendParams)(nil), "apis.SendParams") proto.RegisterType((*UpdateConfigParams)(nil), "apis.UpdateConfigParams") @@ -273,35 +320,38 @@ func init() { proto.RegisterType((*UseridByMobileParams)(nil), "apis.UseridByMobileParams") proto.RegisterType((*Empty)(nil), "apis.Empty") proto.RegisterType((*UseridByMobileReply)(nil), "apis.UseridByMobileReply") + proto.RegisterType((*ValidateConfigReply)(nil), "apis.ValidateConfigReply") } func init() { proto.RegisterFile("send_server.proto", fileDescriptor_63fdd68f7eb311f9) } var fileDescriptor_63fdd68f7eb311f9 = []byte{ - // 357 bytes of a gzipped FileDescriptorProto - 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x74, 0x92, 0xcf, 0x6a, 0xe3, 0x30, - 0x10, 0xc6, 0x71, 0xfe, 0x6e, 0x26, 0x21, 0x64, 0xb5, 0x61, 0xd1, 0xfa, 0x14, 0x0c, 0x59, 0x72, - 0x59, 0x1f, 0xb2, 0x2c, 0x2c, 0xb9, 0x94, 0x36, 0x84, 0x9e, 0x02, 0xc1, 0x4d, 0xce, 0x45, 0x89, - 0xa7, 0x41, 0xd4, 0xb6, 0x8c, 0xa4, 0x04, 0xfc, 0x18, 0x7d, 0x93, 0xd2, 0x27, 0x2c, 0x92, 0xe5, - 0x36, 0x69, 0xd3, 0x9b, 0x7f, 0xdf, 0xe8, 0x1b, 0xcf, 0x37, 0x12, 0x7c, 0x57, 0x98, 0xc5, 0xf7, - 0x0a, 0xe5, 0x11, 0x65, 0x98, 0x4b, 0xa1, 0x05, 0x69, 0xb0, 0x9c, 0xab, 0xe0, 0xd9, 0x03, 0xb8, - 0xc3, 0x2c, 0x5e, 0x31, 0xc9, 0x52, 0x45, 0x28, 0xb4, 0xe7, 0x22, 0xd3, 0x6c, 0xa7, 0xa9, 0x37, - 0xf2, 0x26, 0x9d, 0xa8, 0x42, 0x32, 0x84, 0xe6, 0x5a, 0xe4, 0x7c, 0x47, 0x6b, 0x56, 0x2f, 0xc1, - 0xaa, 0x5c, 0x27, 0x48, 0xeb, 0x4e, 0x35, 0x60, 0xba, 0x2c, 0x51, 0x29, 0xb6, 0x47, 0xda, 0x28, - 0xbb, 0x38, 0x24, 0x3e, 0x7c, 0x5b, 0x49, 0x2e, 0x24, 0xd7, 0x05, 0x6d, 0xda, 0xd2, 0x1b, 0x93, - 0xdf, 0xd0, 0x8f, 0x30, 0x15, 0x1a, 0xd7, 0x98, 0xe6, 0x09, 0xd3, 0x48, 0x5b, 0xf6, 0xc4, 0x07, - 0x35, 0x78, 0xf2, 0x80, 0x6c, 0xf2, 0x98, 0x69, 0x9c, 0x8b, 0xec, 0x81, 0xef, 0xdd, 0xe8, 0x57, - 0xd0, 0xde, 0x59, 0x56, 0xd4, 0x1b, 0xd5, 0x27, 0xdd, 0xe9, 0x38, 0x34, 0x09, 0xc3, 0xcf, 0x47, - 0xc3, 0x12, 0xd4, 0x22, 0xd3, 0xb2, 0x88, 0x2a, 0x97, 0x3f, 0x83, 0xde, 0x69, 0x81, 0x0c, 0xa0, - 0xfe, 0x88, 0x85, 0xdb, 0x83, 0xf9, 0x34, 0x69, 0x8f, 0x2c, 0x39, 0x60, 0xb5, 0x03, 0x0b, 0xb3, - 0xda, 0x7f, 0x2f, 0x08, 0x61, 0xb8, 0x51, 0x28, 0x79, 0x7c, 0x53, 0x2c, 0xc5, 0x96, 0x27, 0xe8, - 0x86, 0xfa, 0x09, 0xad, 0xd4, 0xb2, 0x6b, 0xe3, 0x28, 0x68, 0x43, 0x73, 0x91, 0xe6, 0xba, 0x08, - 0xfe, 0xc0, 0x8f, 0x73, 0x63, 0x84, 0x79, 0x52, 0x18, 0xdf, 0xc1, 0xca, 0x95, 0xaf, 0xa4, 0xe9, - 0x8b, 0x07, 0x1d, 0x73, 0x5d, 0xd7, 0x7b, 0xcc, 0x34, 0x19, 0x43, 0xc3, 0x00, 0x19, 0x94, 0x49, - 0xdf, 0xef, 0xd1, 0xef, 0x96, 0x8a, 0xfd, 0x07, 0xf9, 0x07, 0xbd, 0xd3, 0x25, 0x10, 0xfa, 0xd5, - 0x62, 0xce, 0x6d, 0xb7, 0xd0, 0x3f, 0x1f, 0x8d, 0xf8, 0xce, 0x78, 0x21, 0xa9, 0xff, 0xeb, 0x52, - 0xcd, 0x86, 0xd9, 0xb6, 0xec, 0x83, 0xfb, 0xfb, 0x1a, 0x00, 0x00, 0xff, 0xff, 0xd7, 0x6c, 0xcb, - 0x1b, 0x85, 0x02, 0x00, 0x00, + // 398 bytes of a gzipped FileDescriptorProto + 0x1f, 0x8b, 0x08, 0x00, 0x00, 0x00, 0x00, 0x00, 0x02, 0xff, 0x7c, 0x92, 0xcf, 0x6a, 0xdb, 0x40, + 0x10, 0xc6, 0x51, 0xfc, 0x47, 0xc9, 0x24, 0x98, 0x74, 0x13, 0xca, 0x56, 0xa7, 0x20, 0x70, 0xf1, + 0xa5, 0x3a, 0xb8, 0x14, 0x8a, 0x2f, 0xc5, 0x35, 0xa6, 0x27, 0x83, 0x51, 0xed, 0x5e, 0xcb, 0xda, + 0x9a, 0x8a, 0xa5, 0x92, 0x56, 0xec, 0xae, 0x0d, 0x7a, 0x8c, 0xbe, 0x49, 0x5f, 0xaf, 0xb7, 0xb2, + 0x7f, 0x54, 0x5b, 0xad, 0x9b, 0x9b, 0x7e, 0xdf, 0xcc, 0x37, 0xcc, 0x7c, 0x5a, 0x78, 0xa1, 0xb0, + 0xca, 0xbe, 0x2a, 0x94, 0x47, 0x94, 0x49, 0x2d, 0x85, 0x16, 0xa4, 0xcf, 0x6a, 0xae, 0xe2, 0x9f, + 0x01, 0xc0, 0x67, 0xac, 0xb2, 0x35, 0x93, 0xac, 0x54, 0x84, 0x42, 0xb8, 0x10, 0x95, 0x66, 0x7b, + 0x4d, 0x83, 0xa7, 0x60, 0x72, 0x93, 0xb6, 0x48, 0x1e, 0x61, 0xb0, 0x11, 0x35, 0xdf, 0xd3, 0x2b, + 0xab, 0x3b, 0xb0, 0x2a, 0xd7, 0x05, 0xd2, 0x9e, 0x57, 0x0d, 0x98, 0x29, 0x2b, 0x54, 0x8a, 0xe5, + 0x48, 0xfb, 0x6e, 0x8a, 0x47, 0x12, 0xc1, 0xf5, 0x5a, 0x72, 0x21, 0xb9, 0x6e, 0xe8, 0xc0, 0x96, + 0xfe, 0x30, 0x79, 0x0d, 0xa3, 0x14, 0x4b, 0xa1, 0x71, 0x83, 0x65, 0x5d, 0x30, 0x8d, 0x74, 0x68, + 0x3b, 0xfe, 0x52, 0xe3, 0x1f, 0x01, 0x90, 0x6d, 0x9d, 0x31, 0x8d, 0x0b, 0x51, 0x7d, 0xe3, 0xb9, + 0x5f, 0xfd, 0x03, 0x84, 0x7b, 0xcb, 0x8a, 0x06, 0x4f, 0xbd, 0xc9, 0xed, 0x74, 0x9c, 0x98, 0x0b, + 0x93, 0x7f, 0x5b, 0x13, 0x07, 0x6a, 0x59, 0x69, 0xd9, 0xa4, 0xad, 0x2b, 0x9a, 0xc1, 0xdd, 0x79, + 0x81, 0xdc, 0x43, 0xef, 0x3b, 0x36, 0x3e, 0x07, 0xf3, 0x69, 0xae, 0x3d, 0xb2, 0xe2, 0x80, 0x6d, + 0x06, 0x16, 0x66, 0x57, 0xef, 0x83, 0x38, 0x81, 0xc7, 0xad, 0x42, 0xc9, 0xb3, 0x8f, 0xcd, 0x4a, + 0xec, 0x78, 0x81, 0x7e, 0xa9, 0x97, 0x30, 0x2c, 0x2d, 0xfb, 0x31, 0x9e, 0xe2, 0x10, 0x06, 0xcb, + 0xb2, 0xd6, 0x4d, 0xfc, 0x06, 0x1e, 0xba, 0xc6, 0x14, 0xeb, 0xa2, 0x31, 0xbe, 0x83, 0x95, 0x5b, + 0x9f, 0xa3, 0x78, 0x0e, 0x0f, 0x5f, 0x58, 0xc1, 0x4f, 0x17, 0xb9, 0x76, 0x0a, 0x21, 0x57, 0xb6, + 0x60, 0xfb, 0xaf, 0xd3, 0x16, 0xcd, 0x11, 0xa5, 0xca, 0xfd, 0xc2, 0xe6, 0x73, 0xfa, 0x2b, 0x80, + 0x1b, 0xf3, 0xc7, 0xe7, 0x39, 0x56, 0x9a, 0x8c, 0xa1, 0x6f, 0x80, 0xdc, 0xbb, 0xb0, 0x4e, 0x4f, + 0x21, 0xba, 0x75, 0x8a, 0x5d, 0x93, 0xbc, 0x83, 0xbb, 0xf3, 0x1c, 0x09, 0xfd, 0x5f, 0xb6, 0x5d, + 0xdb, 0x12, 0x46, 0xdd, 0x75, 0x9f, 0x31, 0xbe, 0x72, 0x95, 0x4b, 0xe7, 0x7d, 0x82, 0x51, 0x37, + 0x24, 0x12, 0xf9, 0x31, 0x17, 0x32, 0x6f, 0x07, 0x5d, 0x88, 0x75, 0x37, 0xb4, 0x4f, 0xff, 0xed, + 0xef, 0x00, 0x00, 0x00, 0xff, 0xff, 0xa6, 0x8d, 0x12, 0x77, 0x0f, 0x03, 0x00, 0x00, } // Reference imports to suppress errors if they are not otherwise used. @@ -318,6 +368,7 @@ const _ = grpc.SupportPackageIsVersion4 type SendAgentClient interface { Send(ctx context.Context, in *SendParams, opts ...grpc.CallOption) (*Empty, error) UpdateConfig(ctx context.Context, in *UpdateConfigParams, opts ...grpc.CallOption) (*Empty, error) + ValidateConfig(ctx context.Context, in *UpdateConfigParams, opts ...grpc.CallOption) (*ValidateConfigReply, error) UseridByMobile(ctx context.Context, in *UseridByMobileParams, opts ...grpc.CallOption) (*UseridByMobileReply, error) } @@ -347,6 +398,15 @@ func (c *sendAgentClient) UpdateConfig(ctx context.Context, in *UpdateConfigPara return out, nil } +func (c *sendAgentClient) ValidateConfig(ctx context.Context, in *UpdateConfigParams, opts ...grpc.CallOption) (*ValidateConfigReply, error) { + out := new(ValidateConfigReply) + err := c.cc.Invoke(ctx, "/apis.SendAgent/ValidateConfig", in, out, opts...) + if err != nil { + return nil, err + } + return out, nil +} + func (c *sendAgentClient) UseridByMobile(ctx context.Context, in *UseridByMobileParams, opts ...grpc.CallOption) (*UseridByMobileReply, error) { out := new(UseridByMobileReply) err := c.cc.Invoke(ctx, "/apis.SendAgent/UseridByMobile", in, out, opts...) @@ -360,6 +420,7 @@ func (c *sendAgentClient) UseridByMobile(ctx context.Context, in *UseridByMobile type SendAgentServer interface { Send(context.Context, *SendParams) (*Empty, error) UpdateConfig(context.Context, *UpdateConfigParams) (*Empty, error) + ValidateConfig(context.Context, *UpdateConfigParams) (*ValidateConfigReply, error) UseridByMobile(context.Context, *UseridByMobileParams) (*UseridByMobileReply, error) } @@ -373,6 +434,9 @@ func (*UnimplementedSendAgentServer) Send(ctx context.Context, req *SendParams) func (*UnimplementedSendAgentServer) UpdateConfig(ctx context.Context, req *UpdateConfigParams) (*Empty, error) { return nil, status.Errorf(codes.Unimplemented, "method UpdateConfig not implemented") } +func (*UnimplementedSendAgentServer) ValidateConfig(ctx context.Context, req *UpdateConfigParams) (*ValidateConfigReply, error) { + return nil, status.Errorf(codes.Unimplemented, "method ValidateConfig not implemented") +} func (*UnimplementedSendAgentServer) UseridByMobile(ctx context.Context, req *UseridByMobileParams) (*UseridByMobileReply, error) { return nil, status.Errorf(codes.Unimplemented, "method UseridByMobile not implemented") } @@ -417,6 +481,24 @@ func _SendAgent_UpdateConfig_Handler(srv interface{}, ctx context.Context, dec f return interceptor(ctx, in, info, handler) } +func _SendAgent_ValidateConfig_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(UpdateConfigParams) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(SendAgentServer).ValidateConfig(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: "/apis.SendAgent/ValidateConfig", + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(SendAgentServer).ValidateConfig(ctx, req.(*UpdateConfigParams)) + } + return interceptor(ctx, in, info, handler) +} + func _SendAgent_UseridByMobile_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(UseridByMobileParams) if err := dec(in); err != nil { @@ -447,6 +529,10 @@ var _SendAgent_serviceDesc = grpc.ServiceDesc{ MethodName: "UpdateConfig", Handler: _SendAgent_UpdateConfig_Handler, }, + { + MethodName: "ValidateConfig", + Handler: _SendAgent_ValidateConfig_Handler, + }, { MethodName: "UseridByMobile", Handler: _SendAgent_UseridByMobile_Handler, diff --git a/pkg/notify/rpc/apis/send_server.proto b/pkg/notify/rpc/apis/send_server.proto index c065347486..3652b668ca 100644 --- a/pkg/notify/rpc/apis/send_server.proto +++ b/pkg/notify/rpc/apis/send_server.proto @@ -40,8 +40,14 @@ message UseridByMobileReply { string userid = 1; } +message ValidateConfigReply { + bool isValid = 1; + string msg = 2; +} + service SendAgent { rpc Send(SendParams) returns (Empty); rpc UpdateConfig(UpdateConfigParams) returns (Empty); + rpc ValidateConfig(UpdateConfigParams) returns (ValidateConfigReply); rpc UseridByMobile(UseridByMobileParams) returns (UseridByMobileReply); -} \ No newline at end of file +} diff --git a/pkg/notify/rpc/send.go b/pkg/notify/rpc/send.go index 17a42d892d..c7c4b48c85 100644 --- a/pkg/notify/rpc/send.go +++ b/pkg/notify/rpc/send.go @@ -78,7 +78,7 @@ func (self *SRpcService) InitAll() error { filename := file.Name() if !file.IsDir() && strings.Contains(filename, ".sock") { serviceName := filename[:len(filename)-5] - self.startNewService(ctx, serviceName) + self.startNewService(ctx, serviceName, true) } } if self.SendServices.Len() == 0 { @@ -167,7 +167,7 @@ func (self *SRpcService) execute(ctx context.Context, f func(client *apis.SendNo var err error if !ok { log.Debugf("get service first time failed") - sendService, err = self.startNewService(ctx, serviceName) + sendService, err = self.startNewService(ctx, serviceName, true) if err != nil { return nil, errors.Wrap(err, "start new service failed") @@ -247,7 +247,8 @@ func (self *SRpcService) restartWithConfig(ctx context.Context, serviceName stri } // startNewService try to start a new rpc service named serviceName -func (self *SRpcService) startNewService(ctx context.Context, serviceName string) (*apis.SendNotificationClient, error) { +// passConfig means if pass config to send service +func (self *SRpcService) startNewService(ctx context.Context, serviceName string, passConfig bool) (*apis.SendNotificationClient, error) { var ( sendService *apis.SendNotificationClient @@ -267,6 +268,10 @@ func (self *SRpcService) startNewService(ctx context.Context, serviceName string self.SendServices.Set(sendService, serviceName) + if !passConfig { + return sendService, nil + } + // get config config, err := self.configStore.GetConfig(serviceName) if err != nil { @@ -318,7 +323,7 @@ func (self *SRpcService) updateService(ctx context.Context) error { delete(serviceNameSet, serviceName) continue } - self.startNewService(ctx, serviceName) + self.startNewService(ctx, serviceName, true) } } @@ -331,6 +336,37 @@ func (self *SRpcService) updateService(ctx context.Context) error { return nil } +func (self *SRpcService) ValidateConfig(ctx context.Context, cType string, configs map[string]string) (isValid bool, + message string, err error) { + + sendService, ok := self.SendServices.Get(cType) + + log.Debugf("get service %s", cType) + if !ok { + log.Debugf("get service first time failed") + sendService, err = self.startNewService(ctx, cType, false) + + if err != nil { + err = errors.Wrap(err, "start new service failed") + return + } + } + param := apis.UpdateConfigParams{ + Configs: configs, + } + rep, err := sendService.ValidateConfig(ctx, ¶m) + if err != nil { + st := status.Convert(err) + if st.Code() == codes.Unimplemented { + err = errors.ErrNotImplemented + return + } + err = fmt.Errorf(st.Message()) + return + } + return rep.IsValid, rep.Msg, nil +} + func grpcDialWithUnixSocket(ctx context.Context, socketPath string) (*grpc.ClientConn, error) { return grpc.DialContext(ctx, socketPath, grpc.WithInsecure(), grpc.WithTimeout(time.Second*5), grpc.WithDialer( func(addr string, timeout time.Duration) (net.Conn, error) {