Fix: Message queue trigger update (#2991)
* Handle mqt update * Update scalemanager to support mqtkind update * Short circuit evalution may fail metadata update * Updated mqtmanager test --------- Signed-off-by: Md Soharab Ansari <soharab.ansari@infracloud.io>
This commit is contained in:
@@ -128,7 +128,8 @@ func (opts *UpdateSubCommand) complete(input cli.Input) (err error) {
|
||||
}
|
||||
|
||||
if input.IsSet(flagkey.MqtMetadata) {
|
||||
updated = updated || util.UpdateMapFromStringSlice(&mqt.Spec.Metadata, metadataParams)
|
||||
_ = util.UpdateMapFromStringSlice(&mqt.Spec.Metadata, metadataParams)
|
||||
updated = true
|
||||
}
|
||||
if input.IsSet(flagkey.MqtSecret) {
|
||||
mqt.Spec.Secret = secret
|
||||
|
||||
@@ -42,6 +42,7 @@ const (
|
||||
ADD_TRIGGER requestType = iota
|
||||
DELETE_TRIGGER
|
||||
GET_TRIGGER_SUBSCRIPTION
|
||||
UPDATE_TRIGGER_SUBSCRIPTION
|
||||
)
|
||||
|
||||
type (
|
||||
@@ -166,6 +167,14 @@ func (mqt *MessageQueueTriggerManager) service() {
|
||||
IncreaseSubscriptionCount()
|
||||
}
|
||||
req.respChan <- resp
|
||||
case UPDATE_TRIGGER_SUBSCRIPTION:
|
||||
if _, ok := mqt.triggers[k]; ok {
|
||||
mqt.triggers[k] = req.triggerSub
|
||||
mqt.logger.Debug("updated trigger subscription", zap.String("key", k))
|
||||
} else {
|
||||
resp.err = errors.New("trigger subscription does not exists")
|
||||
}
|
||||
req.respChan <- resp
|
||||
case GET_TRIGGER_SUBSCRIPTION:
|
||||
if _, ok := mqt.triggers[k]; !ok {
|
||||
resp.err = errors.New("trigger does not exist")
|
||||
@@ -198,6 +207,11 @@ func (mqt *MessageQueueTriggerManager) getTriggerSubscription(trigger *fv1.Messa
|
||||
return resp.triggerSub
|
||||
}
|
||||
|
||||
func (mqt *MessageQueueTriggerManager) updateTriggerSubscription(triggerSub *triggerSubscription) error {
|
||||
resp := mqt.makeRequest(UPDATE_TRIGGER_SUBSCRIPTION, triggerSub)
|
||||
return resp.err
|
||||
}
|
||||
|
||||
func (mqt *MessageQueueTriggerManager) checkTriggerSubscription(trigger *fv1.MessageQueueTrigger) bool {
|
||||
return mqt.getTriggerSubscription(trigger) != nil
|
||||
}
|
||||
@@ -207,10 +221,54 @@ func (mqt *MessageQueueTriggerManager) delTriggerSubscription(trigger *fv1.Messa
|
||||
return resp.err
|
||||
}
|
||||
|
||||
func (mqt *MessageQueueTriggerManager) updateTrigger(trigger *fv1.MessageQueueTrigger) error {
|
||||
oldTriggerSubscription := mqt.getTriggerSubscription(trigger)
|
||||
if oldTriggerSubscription == nil {
|
||||
mqt.logger.Info("Trigger subscrption does not exist", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return errors.New("trigger does not exist")
|
||||
}
|
||||
|
||||
// unsubscribe the messagequeue
|
||||
err := mqt.messageQueue.Unsubscribe(oldTriggerSubscription.subscription)
|
||||
if err != nil {
|
||||
mqt.logger.Warn("failed to unsubscribe from message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return err
|
||||
}
|
||||
|
||||
// subscribe using the updated message queue trigger
|
||||
sub, err := mqt.messageQueue.Subscribe(trigger)
|
||||
if err != nil {
|
||||
mqt.logger.Warn("failed to re-subscribe to message queue trigger", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return err
|
||||
}
|
||||
if sub == nil {
|
||||
mqt.logger.Warn("subscription is nil", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return nil
|
||||
}
|
||||
newTriggerSubscription := triggerSubscription{
|
||||
trigger: *trigger,
|
||||
subscription: sub,
|
||||
}
|
||||
|
||||
// update our list
|
||||
err = mqt.updateTriggerSubscription(&newTriggerSubscription)
|
||||
if err != nil {
|
||||
mqt.logger.Fatal("updating message queue trigger failed", zap.Error(err), zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return err
|
||||
}
|
||||
mqt.logger.Info("message queue trigger updated", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
return nil
|
||||
}
|
||||
|
||||
func (mqt *MessageQueueTriggerManager) RegisterTrigger(trigger *fv1.MessageQueueTrigger) error {
|
||||
isPresent := mqt.checkTriggerSubscription(trigger)
|
||||
if isPresent {
|
||||
mqt.logger.Debug("message queue trigger already registered", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
mqt.logger.Debug("updating message queue trigger", zap.String("trigger_name", trigger.ObjectMeta.Name))
|
||||
err := mqt.updateTrigger(trigger)
|
||||
if err != nil {
|
||||
mqt.logger.Error("error updating messagequeuetrigger", zap.Error(err))
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
@@ -30,6 +30,10 @@ import (
|
||||
"github.com/fission/fission/pkg/utils/loggerfactory"
|
||||
)
|
||||
|
||||
const (
|
||||
updatedTopicName = "new-topic"
|
||||
)
|
||||
|
||||
type mqtConsumer struct {
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
@@ -96,7 +100,31 @@ func TestMqtManager(t *testing.T) {
|
||||
if getSub.trigger.ObjectMeta.Name != trigger.ObjectMeta.Name {
|
||||
t.Errorf("getTriggerSubscription should return triggerSub with trigger name %s", trigger.ObjectMeta.Name)
|
||||
}
|
||||
trigger.Spec.Topic = updatedTopicName
|
||||
getSub.subscription.(mqtConsumer).cancel()
|
||||
newSub, err := msgQueue.Subscribe(&trigger)
|
||||
if err != nil {
|
||||
t.Errorf("Subscribe should not return error")
|
||||
}
|
||||
newTriggerSub := triggerSubscription{
|
||||
trigger: trigger,
|
||||
subscription: newSub,
|
||||
}
|
||||
err = mgr.updateTriggerSubscription(&newTriggerSub)
|
||||
if err != nil {
|
||||
t.Errorf("updateTriggerSubscription should not return error")
|
||||
}
|
||||
if !mgr.checkTriggerSubscription(&trigger) {
|
||||
t.Errorf("checkTrigger should return true")
|
||||
}
|
||||
getNewSub := mgr.getTriggerSubscription(&trigger)
|
||||
if getNewSub == nil {
|
||||
t.Fatal("getTriggerSubscription should return triggerSub")
|
||||
}
|
||||
if getNewSub.trigger.Spec.Topic != updatedTopicName {
|
||||
t.Errorf("getTriggerSubscription returns trigger with incorrect topic-name, expected %s got %s", updatedTopicName, getNewSub.trigger.Spec.Topic)
|
||||
}
|
||||
getNewSub.subscription.(mqtConsumer).cancel()
|
||||
err = mgr.delTriggerSubscription(&trigger)
|
||||
if err != nil {
|
||||
t.Errorf("delTriggerSubscription should not return error")
|
||||
|
||||
@@ -41,45 +41,30 @@ func mqTriggerEventHandlers(ctx context.Context, logger *zap.Logger, kubeClient
|
||||
return
|
||||
}
|
||||
logger.Debug("Create deployment for Scaler Object", zap.Any("mqt", mqt.ObjectMeta), zap.Any("mqt.Spec", mqt.Spec))
|
||||
|
||||
authenticationRef := ""
|
||||
if len(mqt.Spec.Secret) > 0 {
|
||||
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
|
||||
err := createAuthTrigger(ctx, kedaClient, mqt, authenticationRef, kubeClient)
|
||||
if err != nil {
|
||||
logger.Error("Failed to create Authentication Trigger", zap.Error(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
if err := createDeployment(ctx, mqt, routerURL, kubeClient); err != nil {
|
||||
logger.Error("Failed to create Deployment", zap.Error(err))
|
||||
if len(authenticationRef) > 0 {
|
||||
err = deleteAuthTrigger(ctx, kedaClient, authenticationRef, mqt.ObjectMeta.Namespace)
|
||||
if err != nil {
|
||||
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if err := createScaledObject(ctx, kedaClient, mqt, authenticationRef); err != nil {
|
||||
logger.Error("Failed to create ScaledObject", zap.Error(err))
|
||||
if len(authenticationRef) > 0 {
|
||||
if err = deleteAuthTrigger(ctx, kedaClient, authenticationRef, mqt.ObjectMeta.Namespace); err != nil {
|
||||
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
|
||||
}
|
||||
}
|
||||
if err = deleteDeployment(ctx, mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace, kubeClient); err != nil {
|
||||
logger.Error("Failed to delete Deployment", zap.Error(err))
|
||||
}
|
||||
}
|
||||
createKedaObjects(ctx, logger, kedaClient, kubeClient, mqt, routerURL)
|
||||
}()
|
||||
},
|
||||
UpdateFunc: func(obj interface{}, newObj interface{}) {
|
||||
go func() {
|
||||
mqt := obj.(*fv1.MessageQueueTrigger)
|
||||
newMqt := newObj.(*fv1.MessageQueueTrigger)
|
||||
mqtkindKedaToFission := (mqt.Spec.MqtKind == "keda" && newMqt.Spec.MqtKind == "fission")
|
||||
mqtkindFissionToKeda := (mqt.Spec.MqtKind == "fission" && newMqt.Spec.MqtKind == "keda")
|
||||
// If mqtkind is updated to fission from keda then
|
||||
// delete keda objects previously created for mqtkind keda.
|
||||
if mqtkindKedaToFission {
|
||||
logger.Debug("Mqtkind updated to fission from keda, cleanup keda objects", zap.Any("mqt", newMqt.ObjectMeta), zap.Any("mqt.Spec", newMqt.Spec))
|
||||
cleanupKedaObjects(ctx, logger, kedaClient, kubeClient, mqt)
|
||||
return
|
||||
}
|
||||
// If mqtkind is updated to keda from fission then
|
||||
// create keda objects
|
||||
if mqtkindFissionToKeda {
|
||||
logger.Debug("Mqtkind changed to keda from fission, create keda objects", zap.Any("mqt", newMqt.ObjectMeta), zap.Any("mqt.Spec", newMqt.Spec))
|
||||
createKedaObjects(ctx, logger, kedaClient, kubeClient, newMqt, routerURL)
|
||||
return
|
||||
}
|
||||
|
||||
updated := checkAndUpdateTriggerFields(mqt, newMqt)
|
||||
if mqt.Spec.MqtKind == "fission" {
|
||||
return
|
||||
@@ -284,6 +269,63 @@ func checkAndUpdateTriggerFields(mqt, newMqt *fv1.MessageQueueTrigger) bool {
|
||||
return updated
|
||||
}
|
||||
|
||||
func createKedaObjects(ctx context.Context, logger *zap.Logger, kedaClient kedaClient.Interface, kubeClient kubernetes.Interface, mqt *fv1.MessageQueueTrigger, routerURL string) {
|
||||
authenticationRef := ""
|
||||
if len(mqt.Spec.Secret) > 0 {
|
||||
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
|
||||
err := createAuthTrigger(ctx, kedaClient, mqt, authenticationRef, kubeClient)
|
||||
if err != nil {
|
||||
logger.Error("Failed to create Authentication Trigger", zap.Error(err))
|
||||
return
|
||||
}
|
||||
}
|
||||
|
||||
if err := createDeployment(ctx, mqt, routerURL, kubeClient); err != nil {
|
||||
logger.Error("Failed to create Deployment", zap.Error(err))
|
||||
if len(authenticationRef) > 0 {
|
||||
err = deleteAuthTrigger(ctx, kedaClient, authenticationRef, mqt.ObjectMeta.Namespace)
|
||||
if err != nil {
|
||||
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
|
||||
}
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
if err := createScaledObject(ctx, kedaClient, mqt, authenticationRef); err != nil {
|
||||
logger.Error("Failed to create ScaledObject", zap.Error(err))
|
||||
if len(authenticationRef) > 0 {
|
||||
if err = deleteAuthTrigger(ctx, kedaClient, authenticationRef, mqt.ObjectMeta.Namespace); err != nil {
|
||||
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
|
||||
}
|
||||
}
|
||||
if err = deleteDeployment(ctx, mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace, kubeClient); err != nil {
|
||||
logger.Error("Failed to delete Deployment", zap.Error(err))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func cleanupKedaObjects(ctx context.Context, logger *zap.Logger, kedaClient kedaClient.Interface, kubeClient kubernetes.Interface, mqt *fv1.MessageQueueTrigger) {
|
||||
authenticationRef := ""
|
||||
if len(mqt.Spec.Secret) > 0 {
|
||||
authenticationRef = fmt.Sprintf("%s-auth-trigger", mqt.ObjectMeta.Name)
|
||||
}
|
||||
|
||||
if len(authenticationRef) > 0 {
|
||||
err := deleteAuthTrigger(ctx, kedaClient, authenticationRef, mqt.ObjectMeta.Namespace)
|
||||
if err != nil {
|
||||
logger.Error("Failed to delete Authentication Trigger", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
if err := deleteDeployment(ctx, mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace, kubeClient); err != nil {
|
||||
logger.Error("Failed to delete Deployment", zap.Error(err))
|
||||
}
|
||||
|
||||
if err := deleteScaledObject(ctx, kedaClient, mqt.ObjectMeta.Name, mqt.ObjectMeta.Namespace); err != nil {
|
||||
logger.Error("Failed to delete ScaledObject", zap.Error(err))
|
||||
}
|
||||
}
|
||||
|
||||
func getAuthTriggerSpec(ctx context.Context, mqt *fv1.MessageQueueTrigger, authenticationRef string, kubeClient kubernetes.Interface) (*kedav1alpha1.TriggerAuthentication, error) {
|
||||
secret, err := kubeClient.CoreV1().Secrets(mqt.Namespace).Get(ctx, mqt.Spec.Secret, metav1.GetOptions{})
|
||||
if err != nil {
|
||||
@@ -517,3 +559,11 @@ func updateScaledObject(ctx context.Context, client kedaClient.Interface, mqt *f
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func deleteScaledObject(ctx context.Context, client kedaClient.Interface, name, namespace string) error {
|
||||
err := client.KedaV1alpha1().ScaledObjects(namespace).Delete(ctx, name, metav1.DeleteOptions{})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user