Commit d7854d72 authored by Kevin Di Lallo's avatar Kevin Di Lallo
Browse files

added app termination confirmation message + moved scenario node removal code...

added app termination confirmation message + moved scenario node removal code from services to sandbox controller
parent 660eb672
Loading
Loading
Loading
Loading
+0 −22
Original line number Diff line number Diff line
@@ -675,28 +675,6 @@ func mec011AppTerminationPost(w http.ResponseWriter, r *http.Request) {
		if sendAppTerminationWhenDone {
			_ = sendTerminationConfirmation(serviceAppInstanceId)
		}

		// Remove node from active scenario
		event := scc.Event{
			Type_: "SCENARIO-UPDATE",
			EventScenarioUpdate: &scc.EventScenarioUpdate{
				Action: "REMOVE",
				Node: &scc.ScenarioNode{
					Type_:  "EDGE-APP",
					Parent: mepName,
					NodeDataUnion: &scc.NodeDataUnion{
						Process: &scc.Process{
							Type_: "EDGE-APP",
							Name:  instanceName,
						},
					},
				},
			},
		}
		_, err := sbxCtrlClient.EventsApi.SendEvent(context.TODO(), event.Type_, event)
		if err != nil {
			log.Error(err)
		}
	}()

	if sendAppTerminationWhenDone {
+20 −3
Original line number Diff line number Diff line
@@ -61,7 +61,7 @@ const (

// MQ payload fields
const (
	mqfieldAppId   string = "id"
	mqFieldAppId   string = "id"
	mqFieldPersist string = "persist"
)

@@ -180,12 +180,12 @@ func msgHandler(msg *mq.Msg, userData interface{}) {
	case mq.MsgAppUpdate:
		log.Debug("RX MSG: ", mq.PrintMsg(msg))
		appStore.Refresh()
		appId := msg.Payload[mqfieldAppId]
		appId := msg.Payload[mqFieldAppId]
		_ = updateAppInfo(appId)
	case mq.MsgAppRemove:
		log.Debug("RX MSG: ", mq.PrintMsg(msg))
		appStore.Refresh()
		appId := msg.Payload[mqfieldAppId]
		appId := msg.Payload[mqFieldAppId]
		_ = terminateAppInfo(appId)
	case mq.MsgAppFlush:
		log.Debug("RX MSG: ", mq.PrintMsg(msg))
@@ -642,6 +642,9 @@ func deleteAppInstance(appId string) {
	// Flush App instance data
	key := baseKey + "app:" + appId
	_ = rc.DBFlush(key)

	// Confirm App removal
	sendAppRemoveCnf(appId)
}

func getAppInfoList() ([]map[string]string, error) {
@@ -916,3 +919,17 @@ func newAppTerminationNotifSubCfg(sub *AppTerminationNotificationSubscription, s
	}
	return subCfg
}

func sendAppRemoveCnf(id string) {
	// Create message to send on MQ
	msg := mqLocal.CreateMsg(mq.MsgAppRemoveCnf, mq.TargetAll, sandboxName)
	msg.Payload[mqFieldAppId] = id

	// Send message to inform other modules of app removal
	log.Debug("TX MSG: ", mq.PrintMsg(msg))
	err := mqLocal.SendMsg(msg)
	if err != nil {
		log.Error("Failed to send message. Error: ", err.Error())
		return
	}
}
+0 −22
Original line number Diff line number Diff line
@@ -3880,28 +3880,6 @@ func mec011AppTerminationPost(w http.ResponseWriter, r *http.Request) {
		if sendAppTerminationWhenDone {
			_ = sendTerminationConfirmation(serviceAppInstanceId)
		}

		// Remove node from active scenario
		event := scc.Event{
			Type_: "SCENARIO-UPDATE",
			EventScenarioUpdate: &scc.EventScenarioUpdate{
				Action: "REMOVE",
				Node: &scc.ScenarioNode{
					Type_:  "EDGE-APP",
					Parent: mepName,
					NodeDataUnion: &scc.NodeDataUnion{
						Process: &scc.Process{
							Type_: "EDGE-APP",
							Name:  instanceName,
						},
					},
				},
			},
		}
		_, err := sbxCtrlClient.EventsApi.SendEvent(context.TODO(), event.Type_, event)
		if err != nil {
			log.Error(err)
		}
	}()

	w.WriteHeader(http.StatusNoContent)
+0 −22
Original line number Diff line number Diff line
@@ -708,28 +708,6 @@ func mec011AppTerminationPost(w http.ResponseWriter, r *http.Request) {
		if sendAppTerminationWhenDone {
			_ = sendTerminationConfirmation(serviceAppInstanceId)
		}

		// Remove node from active scenario
		event := scc.Event{
			Type_: "SCENARIO-UPDATE",
			EventScenarioUpdate: &scc.EventScenarioUpdate{
				Action: "REMOVE",
				Node: &scc.ScenarioNode{
					Type_:  "EDGE-APP",
					Parent: mepName,
					NodeDataUnion: &scc.NodeDataUnion{
						Process: &scc.Process{
							Type_: "EDGE-APP",
							Name:  instanceName,
						},
					},
				},
			},
		}
		_, err := sbxCtrlClient.EventsApi.SendEvent(context.TODO(), event.Type_, event)
		if err != nil {
			log.Error(err)
		}
	}()

	w.WriteHeader(http.StatusNoContent)
+183 −30
Original line number Diff line number Diff line
@@ -34,13 +34,16 @@ import (
)

// MQ payload fields
const mqFieldAppInstanceId = "id"
const mqFieldPersist = "persist"
const (
	mqFieldAppId   = "id"
	mqFieldPersist = "persist"
)

type AppCtrl struct {
	sandboxName string
	appStore    *apps.ApplicationStore
	mqLocal     *mq.MsgQueue
	handlerId   int
}

// App Controller
@@ -74,57 +77,207 @@ func appCtrlInit(sandboxName string, mqLocal *mq.MsgQueue) error {

// Start App Controller
func appCtrlRun() error {
	var err error

	// Register Message Queue handler
	handler := mq.MsgHandler{Handler: msgHandler, UserData: nil}
	appCtrl.handlerId, err = appCtrl.mqLocal.RegisterHandler(handler)
	if err != nil {
		log.Error("Failed to listen for sandbox updates: ", err.Error())
		return err
	}
	return nil
}

// Stop App Controller
func appCtrlStop() error {

	// Unregister handler
	if appCtrl.mqLocal != nil {
		appCtrl.mqLocal.UnregisterHandler(appCtrl.handlerId)
	}
	return nil
}

func appCtrlResetAppInstances(activeModel *mod.Model) error {
	// Flush non-persistent app instances
	appCtrl.appStore.FlushNonPersistent()
// Message Queue handler
func msgHandler(msg *mq.Msg, userData interface{}) {
	switch msg.Message {
	case mq.MsgAppRemoveCnf:
		log.Debug("RX MSG: ", mq.PrintMsg(msg))
		appId := msg.Payload[mqFieldAppId]

	// Create app instances for scenario processes
		// If process exists, remove it from the active scenario
		activeModel := getActiveModel()
		if activeModel != nil {
		// Get active scenario node names
		appNames := activeModel.GetNodeNames(mod.NodeTypeEdgeApp)
		for _, appName := range appNames {
			// Get App Process & context
			appNode := activeModel.GetNode(appName)
			if appNode == nil {
				continue
			proc, ctx, err := getScenarioProcessById(appId, activeModel)
			if err == nil {
				// Prepare scenario update event
				event := &dataModel.Event{
					Type_: "SCENARIO-UPDATE",
					EventScenarioUpdate: &dataModel.EventScenarioUpdate{
						Action: "REMOVE",
						Node: &dataModel.ScenarioNode{
							Type_:  proc.Type_,
							Parent: ctx.Parents[mod.PhyLoc],
							NodeDataUnion: &dataModel.NodeDataUnion{
								Process: proc,
							},
						},
					},
				}
				// Process event to remove node
				_, err = processEvent(event.Type_, event)
				if err != nil {
					log.Error(err.Error())
				}
			}
		}
	default:
	}
			appNodeCtx := activeModel.GetNodeContext(appName)
			if appNodeCtx == nil {
				continue
}
			appProc := appNode.(*dataModel.Process)

func createAppInstance(proc *dataModel.Process, ctx *mod.NodeContext) (*apps.Application, error) {
	// Determine app type
	appType := apps.TypeUser
			if appCtrl.appStore.IsSysApp(appProc.Image) {
	if appCtrl.appStore.IsSysApp(proc.Image) {
		appType = apps.TypeSystem
	}

			// Create & store app instance
	// Create & app instance
	app := &apps.Application{
				Id:      appProc.Id,
				Name:    appProc.Name,
				Mep:     appNodeCtx.Parents[mod.PhyLoc],
		Id:      proc.Id,
		Name:    proc.Name,
		Mep:     ctx.Parents[mod.PhyLoc],
		Type:    appType,
		Persist: false,
	}
			err := appCtrl.appStore.Set(app)
	return app, nil
}

func setAppInstance(name string, activeModel *mod.Model) error {
	// Get scenario Process & context
	proc, ctx, err := getScenarioProcess(name, activeModel)
	if err != nil {
		log.Error(err.Error())
		return err
	}

	// Create app instance
	app, err := createAppInstance(proc, ctx)
	if err != nil {
		log.Error(err.Error())
		return err
	}

	// Set app instance
	err = appCtrl.appStore.Set(app)
	if err != nil {
		log.Error(err.Error())
		return err
	}
	return nil
}

func delAppInstance(id string) error {
	// Validate ID
	if id == "" {
		return errors.New("Invalid app instance ID")
	}

	// Delete app instance
	err := appCtrl.appStore.Del(id)
	if err != nil {
		log.Warn(err.Error())
		return err
	}
	return nil
}

func resetAppInstances(activeModel *mod.Model) error {
	// Flush non-persistent app instances
	appCtrl.appStore.FlushNonPersistent()

	// Get active scenario app list
	scenarioAppList, err := getScenarioAppInstanceList(activeModel)
	if err != nil {
		log.Error(err.Error())
		return err
	}

	// Create app instances for scenario apps
	for _, app := range scenarioAppList {
		err := appCtrl.appStore.Set(app)
		if err != nil {
			log.Error(err.Error())
		}
	}
	return nil
}

func getScenarioAppInstanceList(activeModel *mod.Model) ([]*apps.Application, error) {
	var appList []*apps.Application

	if activeModel != nil {
		// Get active scenario node names
		appNames := activeModel.GetNodeNames(mod.NodeTypeEdgeApp)
		for _, appName := range appNames {
			// Get scenario Process & context
			proc, ctx, err := getScenarioProcess(appName, activeModel)
			if err != nil {
				log.Error(err.Error())
				continue
			}

			// Create app instance
			app, err := createAppInstance(proc, ctx)
			if err != nil {
				log.Error(err.Error())
				continue
			}

			// Add app instance to list
			appList = append(appList, app)
		}
	}
	return appList, nil
}

func getScenarioProcess(name string, activeModel *mod.Model) (*dataModel.Process, *mod.NodeContext, error) {
	// Get app node
	node := activeModel.GetNode(name)
	if node == nil {
		return nil, nil, errors.New("Failed to get app node")
	}
	// Get App Process & context
	proc, ok := node.(*dataModel.Process)
	if !ok {
		return nil, nil, errors.New("Failed to cast node as Process")
	}
	ctx := activeModel.GetNodeContext(proc.Name)
	if ctx == nil {
		return nil, nil, errors.New("Missing node context for " + proc.Name)
	}
	return proc, ctx, nil
}

func getScenarioProcessById(id string, activeModel *mod.Model) (*dataModel.Process, *mod.NodeContext, error) {
	// Get app node
	node := activeModel.GetNodeById(id)
	if node == nil {
		return nil, nil, errors.New("Failed to get app node")
	}
	// Get App Process & context
	proc, ok := node.(*dataModel.Process)
	if !ok {
		return nil, nil, errors.New("Failed to cast node as Process")
	}
	ctx := activeModel.GetNodeContext(proc.Name)
	if ctx == nil {
		return nil, nil, errors.New("Missing node context for " + proc.Name)
	}
	return proc, ctx, nil
}

func applicationsPOST(w http.ResponseWriter, r *http.Request) {
	w.Header().Set("Content-Type", "application/json; charset=UTF-8")
	log.Info("applicationsPOST")
@@ -406,10 +559,10 @@ func appStoreUpdateCb(eventType string, eventData interface{}) {
	switch eventType {
	case apps.EventAdd:
		msg = appCtrl.mqLocal.CreateMsg(mq.MsgAppUpdate, mq.TargetAll, appCtrl.sandboxName)
		msg.Payload[mqFieldAppInstanceId] = eventData.(string)
		msg.Payload[mqFieldAppId] = eventData.(string)
	case apps.EventRemove:
		msg = appCtrl.mqLocal.CreateMsg(mq.MsgAppRemove, mq.TargetAll, appCtrl.sandboxName)
		msg.Payload[mqFieldAppInstanceId] = eventData.(string)
		msg.Payload[mqFieldAppId] = eventData.(string)
	case apps.EventFlush:
		msg = appCtrl.mqLocal.CreateMsg(mq.MsgAppFlush, mq.TargetAll, appCtrl.sandboxName)
		msg.Payload[mqFieldPersist] = strconv.FormatBool(eventData.(bool))
Loading