Loading go-apps/meep-loc-serv/server/loc-serv_test.go +3 −3 Original line number Diff line number Diff line Loading @@ -2042,7 +2042,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone2-poa1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2057,7 +2057,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell2" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2072,7 +2072,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading go-apps/meep-rnis/server/rnis_test.go +3 −3 Original line number Diff line number Diff line Loading @@ -2448,7 +2448,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone2-poa1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2463,7 +2463,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2478,7 +2478,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell2" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading go-apps/meep-sandbox-ctrl/server/sandbox-ctrl.go +39 −11 Original line number Diff line number Diff line Loading @@ -65,6 +65,8 @@ const serviceName = "Sandbox Controller" // MQ payload fields const fieldSandboxName = "sandbox-name" const fieldScenarioName = "scenario-name" const fieldEventType = "event-type" const fieldNodeName = "node-name" // Event types const ( Loading Loading @@ -516,7 +518,7 @@ func ceTerminateScenario(w http.ResponseWriter, r *http.Request) { // Send Terminate message to Virt Engine on Global Message Queue msg = sbxCtrl.mqGlobal.CreateMsg(mq.MsgScenarioTerminate, mq.TargetAll, mq.TargetAll) msg.Payload[fieldSandboxName] = sbxCtrl.sandboxName msg.Payload[fieldScenarioName] = "" msg.Payload[fieldScenarioName] = sbxCtrl.activeModel.GetScenarioName() log.Debug("TX MSG: ", mq.PrintMsg(msg)) err = sbxCtrl.mqGlobal.SendMsg(msg) if err != nil { Loading Loading @@ -634,7 +636,7 @@ func sendEventNetworkCharacteristics(event dataModel.Event) (error, int, string) "throughputUl=" + strconv.Itoa(int(netCharEvent.NetChar.ThroughputUl)) + "Mbps " + "packet-loss=" + strconv.FormatFloat(netCharEvent.NetChar.PacketLoss, 'f', -1, 64) + "% " err := sbxCtrl.activeModel.UpdateNetChar(netCharEvent) err := sbxCtrl.activeModel.UpdateNetChar(netCharEvent, nil) if err != nil { return err, http.StatusInternalServerError, "" } Loading @@ -651,7 +653,7 @@ func sendEventMobility(event dataModel.Event) (error, int, string) { destName := event.EventMobility.Dest description := "[" + elemName + "] move to " + destName oldNL, newNL, err := sbxCtrl.activeModel.MoveNode(elemName, destName) oldNL, newNL, err := sbxCtrl.activeModel.MoveNode(elemName, destName, nil) if err != nil { return err, http.StatusInternalServerError, "" } Loading Loading @@ -737,22 +739,25 @@ func sendEventScenarioUpdate(event dataModel.Event) (error, int, string) { return err, http.StatusBadRequest, "" } // Get node name nodeName := getScenarioNodeName(event.EventScenarioUpdate.Node) // Perform necessary action on scenario switch event.EventScenarioUpdate.Action { case mod.ScenarioAdd: err = sbxCtrl.activeModel.AddScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.AddScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Added node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Added node [" + nodeName + "]" } case mod.ScenarioModify: err = sbxCtrl.activeModel.ModifyScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.ModifyScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Modified node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Modified node [" + nodeName + "]" } case mod.ScenarioRemove: err = sbxCtrl.activeModel.RemoveScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.RemoveScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Removed node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Removed node [" + nodeName + "]" } default: err = errors.New("Unsupported scenario update action: " + event.EventScenarioUpdate.Action) Loading @@ -767,11 +772,16 @@ func sendEventScenarioUpdate(event dataModel.Event) (error, int, string) { // Retrieve element name from type-specific structure func getScenarioNodeName(node *dataModel.ScenarioNode) string { name := "" if node.Type_ == mod.NodeTypeUE { if mod.IsPhyLoc(node.Type_) { if node.NodeDataUnion != nil && node.NodeDataUnion.PhysicalLocation != nil { pl := node.NodeDataUnion.PhysicalLocation name = pl.Name } } else if mod.IsProc(node.Type_) { if node.NodeDataUnion != nil && node.NodeDataUnion.Process != nil { proc := node.NodeDataUnion.Process name = proc.Name } } return name } Loading Loading @@ -1119,10 +1129,28 @@ func ceStopReplayFile(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json; charset=UTF-8") } func activeScenarioUpdateCb() { func activeScenarioUpdateCb(eventType string, userData interface{}) { // Check if update requires Virt Engine intervention if eventType == mod.EventAddNode || eventType == mod.EventRemoveNode || eventType == mod.EventModifyNode { // Send Update message to Virt Engine on Global Message Queue msg := sbxCtrl.mqGlobal.CreateMsg(mq.MsgScenarioUpdate, mq.TargetAll, mq.TargetAll) msg.Payload[fieldSandboxName] = sbxCtrl.sandboxName msg.Payload[fieldEventType] = eventType msg.Payload[fieldNodeName] = userData.(string) log.Debug("TX MSG: ", mq.PrintMsg(msg)) err := sbxCtrl.mqGlobal.SendMsg(msg) if err != nil { log.Error("Failed to send message. Error: ", err.Error()) } } // Send Update message on local Message Queue msg := sbxCtrl.mqLocal.CreateMsg(mq.MsgScenarioUpdate, mq.TargetAll, sbxCtrl.sandboxName) msg.Payload[fieldEventType] = eventType if eventType == mod.EventAddNode || eventType == mod.EventRemoveNode || eventType == mod.EventModifyNode { msg.Payload[fieldNodeName] = userData.(string) } log.Debug("TX MSG: ", mq.PrintMsg(msg)) err := sbxCtrl.mqLocal.SendMsg(msg) if err != nil { Loading go-apps/meep-virt-engine/server/chart-template.go +71 −38 Original line number Diff line number Diff line Loading @@ -200,31 +200,20 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm userChartGroup := strings.Split(proc.UserChartGroup, ":") meSvcName := userChartGroup[1] if meSvcName != "" { if _, found := serviceMap[meSvcName]; !found { serviceMap[meSvcName] = "meepMeSvc: " + meSvcName serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) serviceTemplate.Namespace = scenarioName addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) // NOTE: Every service within a group must expose the same port & protocol var portTemplate ServicePortTemplate portTemplate.Port = userChartGroup[2] portTemplate.Protocol = userChartGroup[3] serviceTemplate.Ports = append(serviceTemplate.Ports, portTemplate) // Create virt-engine chart for new group service chartName := proc.Name + "-svc" chartLocation, err := createChart(chartName, sandboxName, scenarioName, scenarioTemplate) c, err := createMeSvcChart(sandboxName, scenarioName, meSvcName, serviceTemplate.Ports) if err != nil { log.Debug("yaml creation file process: ", err) log.Debug("Failed to create ME Svc chart: ", err) return nil, err } c := newChart(chartName, sandboxName, scenarioName, chartLocation, "") charts = append(charts, c) log.Debug("chart added for user chart group service ", len(charts)) if c != nil { charts = append(charts, *c) log.Debug("chart added for group service: ", meSvcName, " len:", len(charts)) } } } Loading Loading @@ -254,19 +243,7 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) addTemplateLabel(deploymentTemplate, "meepSvc: "+svcName) // Create and store ME Service template only with first occurrence. // If it already exists then add the matching pod label but don't create the service again. meSvcName := proc.ServiceConfig.MeSvcName if meSvcName != "" { if _, found := serviceMap[meSvcName]; !found { serviceMap[meSvcName] = "meepMeSvc: " + meSvcName serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName } addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addTemplateLabel(deploymentTemplate, "meepMeSvc: "+meSvcName) } // Add ports for _, ports := range proc.ServiceConfig.Ports { var portTemplate ServicePortTemplate portTemplate.Port = strconv.Itoa(int(ports.Port)) Loading @@ -283,6 +260,24 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm serviceTemplate.Ports = append(serviceTemplate.Ports, portTemplate) } // Create ME Service chart on first occurrence meSvcName := proc.ServiceConfig.MeSvcName if meSvcName != "" { c, err := createMeSvcChart(sandboxName, scenarioName, meSvcName, serviceTemplate.Ports) if err != nil { log.Debug("Failed to create ME Svc chart: ", err) return nil, err } if c != nil { charts = append(charts, *c) log.Debug("chart added for group service: ", meSvcName, " len:", len(charts)) } // Add ME Svc service & pod labels addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addTemplateLabel(deploymentTemplate, "meepMeSvc: "+meSvcName) } } // Enable GPU template if present Loading Loading @@ -368,6 +363,44 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm return charts, nil } // Create ME Svc chart func createMeSvcChart(sandboxName string, scenarioName string, meSvcName string, ports []ServicePortTemplate) (*helm.Chart, error) { // Ignore if chart already exists if _, found := serviceMap[meSvcName]; found { return nil, nil } // Add ME Svc to map serviceMap[meSvcName] = "meepMeSvc: " + meSvcName // Create default scenario template var scenarioTemplate ScenarioTemplate serviceTemplate := &scenarioTemplate.Service setScenarioDefaults(&scenarioTemplate) // Fill general scenario template information scenarioTemplate.Namespace = scenarioName // Fill ME Svc template information serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName serviceTemplate.Namespace = scenarioName serviceTemplate.Ports = ports addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) // Create virt-engine chart for new group service chartName := "me-svc-" + meSvcName chartLocation, err := createChart(chartName, sandboxName, scenarioName, scenarioTemplate) if err != nil { log.Debug("yaml creation file process: ", err) return nil, err } c := newChart(chartName, sandboxName, scenarioName, chartLocation, "") return &c, nil } func deployCharts(charts []helm.Chart, sandboxName string) error { err := helm.InstallCharts(charts, sandboxName) if err != nil { Loading go-apps/meep-virt-engine/server/virt-engine.go +110 −6 Original line number Diff line number Diff line Loading @@ -23,6 +23,7 @@ import ( "strings" "github.com/InterDigitalInc/AdvantEDGE/go-apps/meep-virt-engine/helm" dataModel "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-data-model" log "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-logger" mod "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-model" mq "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-mq" Loading @@ -38,8 +39,8 @@ const moduleNamespace string = "default" // MQ payload fields const fieldSandboxName = "sandbox-name" // const fieldScenarioName = "scenario-name" const fieldEventType = "event-type" const fieldNodeName = "node-name" type VirtEngine struct { wdPinger *wd.Pinger Loading Loading @@ -175,6 +176,21 @@ func msgHandler(msg *mq.Msg, userData interface{}) { case mq.MsgScenarioActivate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) activateScenario(msg.Payload[fieldSandboxName]) case mq.MsgScenarioUpdate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) eventType := msg.Payload[fieldEventType] sandboxName := msg.Payload[fieldSandboxName] nodeName := msg.Payload[fieldNodeName] switch eventType { case mod.EventAddNode: addScenarioNode(sandboxName, nodeName) case mod.EventModifyNode: modifyScenarioNode(sandboxName, nodeName) case mod.EventRemoveNode: removeScenarioNode(sandboxName, nodeName) default: log.Trace("Ignoring unsupported scenario update event type: ", eventType) } case mq.MsgScenarioTerminate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) terminateScenario(msg.Payload[fieldSandboxName]) Loading Loading @@ -205,6 +221,90 @@ func activateScenario(sandboxName string) { } } func addScenarioNode(sandboxName string, nodeName string) { log.Info("Adding node: ", nodeName) // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Sync with active scenario store activeModel.UpdateScenario() // Find process in active scenario // Create chart template } func modifyScenarioNode(sandboxName string, nodeName string) { log.Info("Modifying node: ", nodeName) // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Sync with active scenario store activeModel.UpdateScenario() } func removeScenarioNode(sandboxName string, nodeName string) { log.Info("Removing node: ", nodeName) if nodeName == "" { log.Error("Missing node name") return } // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Get cached scenario name scenarioName := ve.activeScenarioNames[sandboxName] // Before updating active scenario, find processes to remove procNames := []string{} node := activeModel.GetNode(nodeName) nodeType := activeModel.GetNodeType(nodeName) if mod.IsPhyLoc(nodeType) { pl, ok := node.(*dataModel.PhysicalLocation) if !ok { log.Error("Error casting physical location: " + nodeName) return } for _, proc := range pl.Processes { procNames = append(procNames, proc.Name) } } else if mod.IsProc(nodeType) { proc, ok := node.(*dataModel.Process) if !ok { log.Error("Error casting process: " + nodeName) return } procNames = append(procNames, proc.Name) } else { log.Error("Unsupported node type: ", nodeType) return } // Sync with active scenario store activeModel.UpdateScenario() // Delete releases & remove charts for each process for _, procName := range procNames { _, chartsToDelete := deleteReleases(sandboxName, scenarioName, procName) log.Info("Number of charts to be deleted: ", chartsToDelete) } } func terminateScenario(sandboxName string) { // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] Loading @@ -220,7 +320,7 @@ func terminateScenario(sandboxName string) { scenarioName := ve.activeScenarioNames[sandboxName] // Process right away and start a ticker to retry until everything is deleted _, chartsToDelete := deleteReleases(sandboxName, scenarioName) _, chartsToDelete := deleteReleases(sandboxName, scenarioName, "") log.Info("Number of charts to be deleted: ", chartsToDelete) ve.activeScenarioNames[sandboxName] = "" Loading Loading @@ -269,7 +369,7 @@ func createSandbox(sandboxName string) { func destroySandbox(sandboxName string) { // Process right away and start a ticker to retry until everything is deleted _, chartsToDelete := deleteReleases(sandboxName, "") _, chartsToDelete := deleteReleases(sandboxName, "", "") log.Info("Number of charts to be deleted: ", chartsToDelete) ve.activeScenarioNames[sandboxName] = "" ve.activeModels[sandboxName] = nil Loading @@ -293,7 +393,7 @@ func destroySandbox(sandboxName string) { // }() } func deleteReleases(sandboxName string, scenarioName string) (error, int) { func deleteReleases(sandboxName string, scenarioName string, procName string) (error, int) { if sandboxName == "" { return nil, 0 } Loading @@ -302,9 +402,13 @@ func deleteReleases(sandboxName string, scenarioName string) (error, int) { path := "/charts/" + sandboxName releasePrefix := "meep-" if scenarioName != "" { path += "/scenario/" path += "/scenario/" + scenarioName releasePrefix += scenarioName + "-" } if procName != "" { path += "/" + procName releasePrefix += procName } // Retrieve list of releases chartsToDelete := 0 Loading Loading
go-apps/meep-loc-serv/server/loc-serv_test.go +3 −3 Original line number Diff line number Diff line Loading @@ -2042,7 +2042,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone2-poa1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2057,7 +2057,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell2" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2072,7 +2072,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading
go-apps/meep-rnis/server/rnis_test.go +3 −3 Original line number Diff line number Diff line Loading @@ -2448,7 +2448,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone2-poa1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2463,7 +2463,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell1" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading @@ -2478,7 +2478,7 @@ func updateScenario(testUpdate string) { elemName := "ue1" destName := "zone1-poa-cell2" _, _, err := m.MoveNode(elemName, destName) _, _, err := m.MoveNode(elemName, destName, nil) if err != nil { log.Error("Error sending mobility event") } Loading
go-apps/meep-sandbox-ctrl/server/sandbox-ctrl.go +39 −11 Original line number Diff line number Diff line Loading @@ -65,6 +65,8 @@ const serviceName = "Sandbox Controller" // MQ payload fields const fieldSandboxName = "sandbox-name" const fieldScenarioName = "scenario-name" const fieldEventType = "event-type" const fieldNodeName = "node-name" // Event types const ( Loading Loading @@ -516,7 +518,7 @@ func ceTerminateScenario(w http.ResponseWriter, r *http.Request) { // Send Terminate message to Virt Engine on Global Message Queue msg = sbxCtrl.mqGlobal.CreateMsg(mq.MsgScenarioTerminate, mq.TargetAll, mq.TargetAll) msg.Payload[fieldSandboxName] = sbxCtrl.sandboxName msg.Payload[fieldScenarioName] = "" msg.Payload[fieldScenarioName] = sbxCtrl.activeModel.GetScenarioName() log.Debug("TX MSG: ", mq.PrintMsg(msg)) err = sbxCtrl.mqGlobal.SendMsg(msg) if err != nil { Loading Loading @@ -634,7 +636,7 @@ func sendEventNetworkCharacteristics(event dataModel.Event) (error, int, string) "throughputUl=" + strconv.Itoa(int(netCharEvent.NetChar.ThroughputUl)) + "Mbps " + "packet-loss=" + strconv.FormatFloat(netCharEvent.NetChar.PacketLoss, 'f', -1, 64) + "% " err := sbxCtrl.activeModel.UpdateNetChar(netCharEvent) err := sbxCtrl.activeModel.UpdateNetChar(netCharEvent, nil) if err != nil { return err, http.StatusInternalServerError, "" } Loading @@ -651,7 +653,7 @@ func sendEventMobility(event dataModel.Event) (error, int, string) { destName := event.EventMobility.Dest description := "[" + elemName + "] move to " + destName oldNL, newNL, err := sbxCtrl.activeModel.MoveNode(elemName, destName) oldNL, newNL, err := sbxCtrl.activeModel.MoveNode(elemName, destName, nil) if err != nil { return err, http.StatusInternalServerError, "" } Loading Loading @@ -737,22 +739,25 @@ func sendEventScenarioUpdate(event dataModel.Event) (error, int, string) { return err, http.StatusBadRequest, "" } // Get node name nodeName := getScenarioNodeName(event.EventScenarioUpdate.Node) // Perform necessary action on scenario switch event.EventScenarioUpdate.Action { case mod.ScenarioAdd: err = sbxCtrl.activeModel.AddScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.AddScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Added node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Added node [" + nodeName + "]" } case mod.ScenarioModify: err = sbxCtrl.activeModel.ModifyScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.ModifyScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Modified node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Modified node [" + nodeName + "]" } case mod.ScenarioRemove: err = sbxCtrl.activeModel.RemoveScenarioNode(event.EventScenarioUpdate.Node) err = sbxCtrl.activeModel.RemoveScenarioNode(event.EventScenarioUpdate.Node, nodeName) if err == nil { description = "Removed node [" + getScenarioNodeName(event.EventScenarioUpdate.Node) + "]" description = "Removed node [" + nodeName + "]" } default: err = errors.New("Unsupported scenario update action: " + event.EventScenarioUpdate.Action) Loading @@ -767,11 +772,16 @@ func sendEventScenarioUpdate(event dataModel.Event) (error, int, string) { // Retrieve element name from type-specific structure func getScenarioNodeName(node *dataModel.ScenarioNode) string { name := "" if node.Type_ == mod.NodeTypeUE { if mod.IsPhyLoc(node.Type_) { if node.NodeDataUnion != nil && node.NodeDataUnion.PhysicalLocation != nil { pl := node.NodeDataUnion.PhysicalLocation name = pl.Name } } else if mod.IsProc(node.Type_) { if node.NodeDataUnion != nil && node.NodeDataUnion.Process != nil { proc := node.NodeDataUnion.Process name = proc.Name } } return name } Loading Loading @@ -1119,10 +1129,28 @@ func ceStopReplayFile(w http.ResponseWriter, r *http.Request) { w.Header().Set("Content-Type", "application/json; charset=UTF-8") } func activeScenarioUpdateCb() { func activeScenarioUpdateCb(eventType string, userData interface{}) { // Check if update requires Virt Engine intervention if eventType == mod.EventAddNode || eventType == mod.EventRemoveNode || eventType == mod.EventModifyNode { // Send Update message to Virt Engine on Global Message Queue msg := sbxCtrl.mqGlobal.CreateMsg(mq.MsgScenarioUpdate, mq.TargetAll, mq.TargetAll) msg.Payload[fieldSandboxName] = sbxCtrl.sandboxName msg.Payload[fieldEventType] = eventType msg.Payload[fieldNodeName] = userData.(string) log.Debug("TX MSG: ", mq.PrintMsg(msg)) err := sbxCtrl.mqGlobal.SendMsg(msg) if err != nil { log.Error("Failed to send message. Error: ", err.Error()) } } // Send Update message on local Message Queue msg := sbxCtrl.mqLocal.CreateMsg(mq.MsgScenarioUpdate, mq.TargetAll, sbxCtrl.sandboxName) msg.Payload[fieldEventType] = eventType if eventType == mod.EventAddNode || eventType == mod.EventRemoveNode || eventType == mod.EventModifyNode { msg.Payload[fieldNodeName] = userData.(string) } log.Debug("TX MSG: ", mq.PrintMsg(msg)) err := sbxCtrl.mqLocal.SendMsg(msg) if err != nil { Loading
go-apps/meep-virt-engine/server/chart-template.go +71 −38 Original line number Diff line number Diff line Loading @@ -200,31 +200,20 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm userChartGroup := strings.Split(proc.UserChartGroup, ":") meSvcName := userChartGroup[1] if meSvcName != "" { if _, found := serviceMap[meSvcName]; !found { serviceMap[meSvcName] = "meepMeSvc: " + meSvcName serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) serviceTemplate.Namespace = scenarioName addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) // NOTE: Every service within a group must expose the same port & protocol var portTemplate ServicePortTemplate portTemplate.Port = userChartGroup[2] portTemplate.Protocol = userChartGroup[3] serviceTemplate.Ports = append(serviceTemplate.Ports, portTemplate) // Create virt-engine chart for new group service chartName := proc.Name + "-svc" chartLocation, err := createChart(chartName, sandboxName, scenarioName, scenarioTemplate) c, err := createMeSvcChart(sandboxName, scenarioName, meSvcName, serviceTemplate.Ports) if err != nil { log.Debug("yaml creation file process: ", err) log.Debug("Failed to create ME Svc chart: ", err) return nil, err } c := newChart(chartName, sandboxName, scenarioName, chartLocation, "") charts = append(charts, c) log.Debug("chart added for user chart group service ", len(charts)) if c != nil { charts = append(charts, *c) log.Debug("chart added for group service: ", meSvcName, " len:", len(charts)) } } } Loading Loading @@ -254,19 +243,7 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) addTemplateLabel(deploymentTemplate, "meepSvc: "+svcName) // Create and store ME Service template only with first occurrence. // If it already exists then add the matching pod label but don't create the service again. meSvcName := proc.ServiceConfig.MeSvcName if meSvcName != "" { if _, found := serviceMap[meSvcName]; !found { serviceMap[meSvcName] = "meepMeSvc: " + meSvcName serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName } addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addTemplateLabel(deploymentTemplate, "meepMeSvc: "+meSvcName) } // Add ports for _, ports := range proc.ServiceConfig.Ports { var portTemplate ServicePortTemplate portTemplate.Port = strconv.Itoa(int(ports.Port)) Loading @@ -283,6 +260,24 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm serviceTemplate.Ports = append(serviceTemplate.Ports, portTemplate) } // Create ME Service chart on first occurrence meSvcName := proc.ServiceConfig.MeSvcName if meSvcName != "" { c, err := createMeSvcChart(sandboxName, scenarioName, meSvcName, serviceTemplate.Ports) if err != nil { log.Debug("Failed to create ME Svc chart: ", err) return nil, err } if c != nil { charts = append(charts, *c) log.Debug("chart added for group service: ", meSvcName, " len:", len(charts)) } // Add ME Svc service & pod labels addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addTemplateLabel(deploymentTemplate, "meepMeSvc: "+meSvcName) } } // Enable GPU template if present Loading Loading @@ -368,6 +363,44 @@ func generateScenarioCharts(sandboxName string, model *mod.Model) (charts []helm return charts, nil } // Create ME Svc chart func createMeSvcChart(sandboxName string, scenarioName string, meSvcName string, ports []ServicePortTemplate) (*helm.Chart, error) { // Ignore if chart already exists if _, found := serviceMap[meSvcName]; found { return nil, nil } // Add ME Svc to map serviceMap[meSvcName] = "meepMeSvc: " + meSvcName // Create default scenario template var scenarioTemplate ScenarioTemplate serviceTemplate := &scenarioTemplate.Service setScenarioDefaults(&scenarioTemplate) // Fill general scenario template information scenarioTemplate.Namespace = scenarioName // Fill ME Svc template information serviceTemplate.MeServiceEnabled = trueStr serviceTemplate.MeServiceName = meSvcName serviceTemplate.Namespace = scenarioName serviceTemplate.Ports = ports addServiceLabel(serviceTemplate, "meepMeSvc: "+meSvcName) addServiceLabel(serviceTemplate, "meepScenario: "+scenarioName) // Create virt-engine chart for new group service chartName := "me-svc-" + meSvcName chartLocation, err := createChart(chartName, sandboxName, scenarioName, scenarioTemplate) if err != nil { log.Debug("yaml creation file process: ", err) return nil, err } c := newChart(chartName, sandboxName, scenarioName, chartLocation, "") return &c, nil } func deployCharts(charts []helm.Chart, sandboxName string) error { err := helm.InstallCharts(charts, sandboxName) if err != nil { Loading
go-apps/meep-virt-engine/server/virt-engine.go +110 −6 Original line number Diff line number Diff line Loading @@ -23,6 +23,7 @@ import ( "strings" "github.com/InterDigitalInc/AdvantEDGE/go-apps/meep-virt-engine/helm" dataModel "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-data-model" log "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-logger" mod "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-model" mq "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-mq" Loading @@ -38,8 +39,8 @@ const moduleNamespace string = "default" // MQ payload fields const fieldSandboxName = "sandbox-name" // const fieldScenarioName = "scenario-name" const fieldEventType = "event-type" const fieldNodeName = "node-name" type VirtEngine struct { wdPinger *wd.Pinger Loading Loading @@ -175,6 +176,21 @@ func msgHandler(msg *mq.Msg, userData interface{}) { case mq.MsgScenarioActivate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) activateScenario(msg.Payload[fieldSandboxName]) case mq.MsgScenarioUpdate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) eventType := msg.Payload[fieldEventType] sandboxName := msg.Payload[fieldSandboxName] nodeName := msg.Payload[fieldNodeName] switch eventType { case mod.EventAddNode: addScenarioNode(sandboxName, nodeName) case mod.EventModifyNode: modifyScenarioNode(sandboxName, nodeName) case mod.EventRemoveNode: removeScenarioNode(sandboxName, nodeName) default: log.Trace("Ignoring unsupported scenario update event type: ", eventType) } case mq.MsgScenarioTerminate: log.Debug("RX MSG: ", mq.PrintMsg(msg)) terminateScenario(msg.Payload[fieldSandboxName]) Loading Loading @@ -205,6 +221,90 @@ func activateScenario(sandboxName string) { } } func addScenarioNode(sandboxName string, nodeName string) { log.Info("Adding node: ", nodeName) // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Sync with active scenario store activeModel.UpdateScenario() // Find process in active scenario // Create chart template } func modifyScenarioNode(sandboxName string, nodeName string) { log.Info("Modifying node: ", nodeName) // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Sync with active scenario store activeModel.UpdateScenario() } func removeScenarioNode(sandboxName string, nodeName string) { log.Info("Removing node: ", nodeName) if nodeName == "" { log.Error("Missing node name") return } // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] if activeModel == nil { log.Error("No active model for sandbox: ", sandboxName) return } // Get cached scenario name scenarioName := ve.activeScenarioNames[sandboxName] // Before updating active scenario, find processes to remove procNames := []string{} node := activeModel.GetNode(nodeName) nodeType := activeModel.GetNodeType(nodeName) if mod.IsPhyLoc(nodeType) { pl, ok := node.(*dataModel.PhysicalLocation) if !ok { log.Error("Error casting physical location: " + nodeName) return } for _, proc := range pl.Processes { procNames = append(procNames, proc.Name) } } else if mod.IsProc(nodeType) { proc, ok := node.(*dataModel.Process) if !ok { log.Error("Error casting process: " + nodeName) return } procNames = append(procNames, proc.Name) } else { log.Error("Unsupported node type: ", nodeType) return } // Sync with active scenario store activeModel.UpdateScenario() // Delete releases & remove charts for each process for _, procName := range procNames { _, chartsToDelete := deleteReleases(sandboxName, scenarioName, procName) log.Info("Number of charts to be deleted: ", chartsToDelete) } } func terminateScenario(sandboxName string) { // Get sandbox-specific active model activeModel := ve.activeModels[sandboxName] Loading @@ -220,7 +320,7 @@ func terminateScenario(sandboxName string) { scenarioName := ve.activeScenarioNames[sandboxName] // Process right away and start a ticker to retry until everything is deleted _, chartsToDelete := deleteReleases(sandboxName, scenarioName) _, chartsToDelete := deleteReleases(sandboxName, scenarioName, "") log.Info("Number of charts to be deleted: ", chartsToDelete) ve.activeScenarioNames[sandboxName] = "" Loading Loading @@ -269,7 +369,7 @@ func createSandbox(sandboxName string) { func destroySandbox(sandboxName string) { // Process right away and start a ticker to retry until everything is deleted _, chartsToDelete := deleteReleases(sandboxName, "") _, chartsToDelete := deleteReleases(sandboxName, "", "") log.Info("Number of charts to be deleted: ", chartsToDelete) ve.activeScenarioNames[sandboxName] = "" ve.activeModels[sandboxName] = nil Loading @@ -293,7 +393,7 @@ func destroySandbox(sandboxName string) { // }() } func deleteReleases(sandboxName string, scenarioName string) (error, int) { func deleteReleases(sandboxName string, scenarioName string, procName string) (error, int) { if sandboxName == "" { return nil, 0 } Loading @@ -302,9 +402,13 @@ func deleteReleases(sandboxName string, scenarioName string) (error, int) { path := "/charts/" + sandboxName releasePrefix := "meep-" if scenarioName != "" { path += "/scenario/" path += "/scenario/" + scenarioName releasePrefix += scenarioName + "-" } if procName != "" { path += "/" + procName releasePrefix += procName } // Retrieve list of releases chartsToDelete := 0 Loading