mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
feature(notify): Validate config
This commit is contained in:
@@ -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
|
||||
})
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
+65
-12
@@ -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["<type>"]
|
||||
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.
|
||||
|
||||
@@ -139,6 +139,9 @@ func AddNotifyDispatcher(prefix string, app *appsrv.Application) {
|
||||
app.AddHandler2("DELETE",
|
||||
fmt.Sprintf("%s/%s/<type>", prefix, modelDispatcher.KeywordPlural()),
|
||||
middleware(configDeleteHandler), metadata, "delete_configs", tags)
|
||||
app.AddHandler2("POST",
|
||||
fmt.Sprintf("%s/%s/<type>/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["<type>"]
|
||||
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())
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
|
||||
+40
-4
@@ -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) {
|
||||
|
||||
Reference in New Issue
Block a user