diff --git a/pkg/fission-cli/cmd/mqtrigger/update.go b/pkg/fission-cli/cmd/mqtrigger/update.go index 14dbe8e8..e40f8af1 100644 --- a/pkg/fission-cli/cmd/mqtrigger/update.go +++ b/pkg/fission-cli/cmd/mqtrigger/update.go @@ -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 diff --git a/pkg/mqtrigger/mqtmanager.go b/pkg/mqtrigger/mqtmanager.go index 619f9afe..a0ab1e0e 100644 --- a/pkg/mqtrigger/mqtmanager.go +++ b/pkg/mqtrigger/mqtmanager.go @@ -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 } diff --git a/pkg/mqtrigger/mqtmanager_test.go b/pkg/mqtrigger/mqtmanager_test.go index d6d93969..aad79d14 100644 --- a/pkg/mqtrigger/mqtmanager_test.go +++ b/pkg/mqtrigger/mqtmanager_test.go @@ -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") diff --git a/pkg/mqtrigger/scalermanager.go b/pkg/mqtrigger/scalermanager.go index 58cac40b..b9846246 100644 --- a/pkg/mqtrigger/scalermanager.go +++ b/pkg/mqtrigger/scalermanager.go @@ -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 +}