Commit 9882d8db authored by Muhammad Umair Khan's avatar Muhammad Umair Khan
Browse files

Bug fix: TC-Engine

parent ff9a85a1
Loading
Loading
Loading
Loading
+3 −1
Original line number Diff line number Diff line
@@ -16,4 +16,6 @@ data:
    - name: init-{{ .Values.sidecar.dependency }}
      image: busybox:1.28
      imagePullPolicy: IfNotPresent
      command: ['sh', '-c', 'until nslookup {{ .Values.sidecar.dependency }}.kube-system ; do echo waiting for {{ .Values.sidecar.dependency }}; sleep 0.25; done;']
      securityContext:
        privileged: true
      command: ['sh', '-c', 'sysctl -w net.ipv4.ip_forward=1 || true; until nslookup {{ .Values.sidecar.dependency }}.kube-system ; do echo waiting for {{ .Values.sidecar.dependency }}; sleep 0.25; done;']
+93 −7
Original line number Diff line number Diff line
@@ -58,9 +58,11 @@ const svcPrefix string = "SVC-"
const mePrefix string = meepPrefix + "ME-"
const ingressPrefix string = meepPrefix + "INGRESS-"
const egressPrefix string = meepPrefix + "EGRESS-"
const egressSnatPrefix string = meepPrefix + "E-SNAT-"
const meSvcChain string = mePrefix + "SERVICES"
const ingressSvcChain string = ingressPrefix + "SERVICES"
const egressSvcChain string = egressPrefix + "SERVICES"
const egressSnatChain string = egressSnatPrefix + "SERVICES"
const maxChainLen int = 25
const capLetters string = "ABCDEFGHIJKLMNOPQRSTUVWXYZ"
const ipAddrNone string = "n/a"
@@ -358,13 +360,6 @@ func refreshLbRules() {
		}
	}

	// Reapply masquerading rule if not present
	err = ipTbl.AppendUnique("nat", "POSTROUTING", "-o", "eth0", "-j", "MASQUERADE")
	if err != nil {
		log.Error("Failed to set rule [-A POSTROUTING -o eth0 -j MASQUERADE]. Error: ", err)
		return
	}

	// Create top-level MEEP service chains if not present
	// MEEP-ME-SERVICES
	_, exists := chainMap[meSvcChain]
@@ -402,6 +397,18 @@ func refreshLbRules() {
	}
	delete(chainMap, egressSvcChain)

	// MEEP-E-SNAT-SERVICES
	_, exists = chainMap[egressSnatChain]
	if !exists {
		log.Debug("Creating MEEP chain MEEP-E-SNAT-SERVICES")
		err = ipTbl.NewChain("nat", egressSnatChain)
		if err != nil {
			log.Error("Failed to create chain. Error: ", err)
			return
		}
	}
	delete(chainMap, egressSnatChain)

	// Reapply top-level routing rules if not present
	err = ipTbl.AppendUnique("nat", "OUTPUT", "-j", meSvcChain)
	if err != nil {
@@ -418,6 +425,11 @@ func refreshLbRules() {
		log.Error("Failed to set rule [-A PREROUTING -j "+egressSvcChain+"]. Error: ", err)
		return
	}
	err = ipTbl.AppendUnique("nat", "POSTROUTING", "-o", "eth0", "-j", egressSnatChain)
	if err != nil {
		log.Error("Failed to set rule [-A POSTROUTING -o eth0 -j "+egressSnatChain+"]. Error: ", err)
		return
	}

	// Apply pod-specific LB rules stored in DB
	flushRequired = false
@@ -435,6 +447,8 @@ func refreshLbRules() {

		if strings.Contains(chain, ingressPrefix) {
			parentChain = ingressSvcChain
		} else if strings.Contains(chain, egressSnatPrefix) {
			parentChain = egressSnatChain
		} else if strings.Contains(chain, egressPrefix) {
			parentChain = egressSvcChain
		} else {
@@ -543,6 +557,12 @@ func refreshLbRulesHandler(key string, fields map[string]string, userData interf
			// No update required. Remove chain from chain map and return.
			if exists {
				delete(*chainMap, serviceChain)
				if fields[fieldSvcType] == typeEgressSvc {
					err = addEgressSnatRule(fields, chainMap)
					if err != nil {
						return err
					}
				}
				return nil
			}
		}
@@ -574,6 +594,72 @@ func refreshLbRulesHandler(key string, fields map[string]string, userData interf
		return err
	}

	// For Egress services, also create destination-scoped SNAT/MASQUERADE rule in MEEP-E-SNAT-SERVICES
	if fields[fieldSvcType] == typeEgressSvc {
		err = addEgressSnatRule(fields, chainMap)
		if err != nil {
			return err
		}
	}

	flushRequired = true
	return nil
}

func addEgressSnatRule(fields map[string]string, chainMap *map[string]bool) error {
	var err error
	servicePrefix := egressSnatPrefix + svcPrefix
	service := servicePrefix + strings.ToUpper(fields[fieldSvcName]) + "-" + fields[fieldSvcPort]
	var args []string
	args = append(args, "-p", fields[fieldSvcProtocol], "-d", fields[fieldLbSvcIp], "--dport", fields[fieldLbSvcPort],
		"-j", "MASQUERADE", "-m", "comment", "--comment", service)

	// Retrieve service chain name if service exists
	serviceChain, exists := serviceChains[service]
	if exists {
		// Check if chain exists
		_, exists = (*chainMap)[serviceChain]
		if exists {
			// Check if rule requires update
			exists, err = ipTbl.Exists("nat", serviceChain, args...)
			if err != nil {
				log.Error("Failed to check if rule exists. Error: ", err)
				return err
			}

			// No update required. Remove chain from chain map and return.
			if exists {
				delete(*chainMap, serviceChain)
				return nil
			}
		}
	}

	// Create new service chain name
	log.Debug("Creating new service chain mapping for SNAT service: ", service)
	serviceChain = servicePrefix + randSeq(maxChainLen-len(servicePrefix))
	serviceChains[service] = serviceChain

	// Create MEEP service chain
	log.Debug("Creating MEEP chain ", serviceChain)
	err = ipTbl.NewChain("nat", serviceChain)
	if err != nil {
		log.Error("Failed to create chain. Error: ", err)
		return err
	}

	// Create service routing rules
	err = ipTbl.AppendUnique("nat", egressSnatChain, "-j", serviceChain)
	if err != nil {
		log.Error("Failed to set rule [-A ", egressSnatChain, " -j ", serviceChain, "]. Error: ", err)
		return err
	}
	err = ipTbl.AppendUnique("nat", serviceChain, args...)
	if err != nil {
		log.Error("Failed to set rule [-A ", egressSnatChain, " -j ", serviceChain, " ", args, "]. Error: ", err)
		return err
	}

	flushRequired = true
	return nil
}
+4 −0
Original line number Diff line number Diff line
@@ -36,3 +36,7 @@ func InstallCharts(charts []Chart, sandboxName string) error {
func DeleteReleases(charts []Chart, sandboxName string) error {
	return runTask(Delete, charts, sandboxName)
}

func CleanChartDir(chartDir string) error {
	return runCleanTask(chartDir)
}
+18 −0
Original line number Diff line number Diff line
@@ -17,6 +17,8 @@
package helm

import (
	"os"

	log "github.com/InterDigitalInc/AdvantEDGE/go-packages/meep-logger"
)

@@ -25,12 +27,14 @@ type Task string
const (
	Install Task = "INSTALL"
	Delete  Task = "DELETE"
	Clean   Task = "CLEAN"
)

type Job struct {
	task        Task
	charts      []Chart
	sandboxName string
	chartDir    string
}

var queue *chan Job = nil
@@ -54,6 +58,13 @@ func startWorker() {
				log.Debug("Deleting ", len(job.charts), " Releases...")
				_ = deleteReleases(job.charts)
				log.Debug("Releases deleted (", len(job.charts), ")")

			case Clean:
				log.Debug("Removing chart directory: ", job.chartDir)
				if _, err := os.Stat(job.chartDir); err == nil {
					_ = os.RemoveAll(job.chartDir)
				}
				log.Debug("Chart directory removed (", job.chartDir, ")")
			}
		}
		queue = nil
@@ -66,3 +77,10 @@ func runTask(task Task, charts []Chart, sandboxName string) error {
	*queue <- job
	return nil
}

func runCleanTask(chartDir string) error {
	startWorker()
	var job Job = Job{task: Clean, chartDir: chartDir}
	*queue <- job
	return nil
}
+59 −15
Original line number Diff line number Diff line
@@ -211,9 +211,30 @@ func msgHandler(msg *mq.Msg, userData interface{}) {
	}
}

func getModel(sandboxName string) *mod.Model {
	activeModel := ve.activeModels[sandboxName]
	if activeModel == nil {
		modelCfg := mod.ModelCfg{
			Name:      moduleName,
			Namespace: sandboxName,
			Module:    moduleName,
			DbAddr:    redisAddr,
			UpdateCb:  nil,
		}
		var err error
		activeModel, err = mod.NewModel(modelCfg)
		if err != nil {
			log.Error("Failed to create model: ", err.Error())
			return nil
		}
		ve.activeModels[sandboxName] = activeModel
	}
	return activeModel
}

func activateScenario(sandboxName string) {
	// Get sandbox-specific active model
	activeModel := ve.activeModels[sandboxName]
	activeModel := getModel(sandboxName)
	if activeModel == nil {
		log.Error("No active model for sandbox: ", sandboxName)
		return
@@ -237,7 +258,7 @@ func addScenarioNode(sandboxName string, nodeName string) {
	log.Info("Adding node: ", nodeName)

	// Get sandbox-specific active model
	activeModel := ve.activeModels[sandboxName]
	activeModel := getModel(sandboxName)
	if activeModel == nil {
		log.Error("No active model for sandbox: ", sandboxName)
		return
@@ -264,7 +285,7 @@ func modifyScenarioNode(sandboxName string, nodeName string) {
	log.Info("Modifying node: ", nodeName)

	// Get sandbox-specific active model
	activeModel := ve.activeModels[sandboxName]
	activeModel := getModel(sandboxName)
	if activeModel == nil {
		log.Error("No active model for sandbox: ", sandboxName)
		return
@@ -272,6 +293,10 @@ func modifyScenarioNode(sandboxName string, nodeName string) {

	// Get cached scenario name
	scenarioName := ve.activeScenarioNames[sandboxName]
	if scenarioName == "" {
		scenarioName = activeModel.GetScenarioName()
		ve.activeScenarioNames[sandboxName] = scenarioName
	}

	// Sync with active scenario store
	activeModel.UpdateScenario()
@@ -304,7 +329,7 @@ func removeScenarioNode(sandboxName string, nodeName string) {
	}

	// Get sandbox-specific active model
	activeModel := ve.activeModels[sandboxName]
	activeModel := getModel(sandboxName)
	if activeModel == nil {
		log.Error("No active model for sandbox: ", sandboxName)
		return
@@ -312,6 +337,10 @@ func removeScenarioNode(sandboxName string, nodeName string) {

	// Get cached scenario name
	scenarioName := ve.activeScenarioNames[sandboxName]
	if scenarioName == "" {
		scenarioName = activeModel.GetScenarioName()
		ve.activeScenarioNames[sandboxName] = scenarioName
	}

	// Before updating active scenario, find processes to remove
	procNames := []string{}
@@ -358,6 +387,9 @@ func terminateScenario(sandboxName string, scenarioName string) {
	if scenarioName == "" {
		// Get cached scenario name
		scenarioName = ve.activeScenarioNames[sandboxName]
		if scenarioName == "" && ve.activeModels[sandboxName] != nil {
			scenarioName = ve.activeModels[sandboxName].GetScenarioName()
		}
	}

	if scenarioName == "" {
@@ -370,8 +402,8 @@ func terminateScenario(sandboxName string, scenarioName string) {
	log.Info("Number of charts to be deleted: ", chartsToDelete)
	ve.activeScenarioNames[sandboxName] = ""

	// Clean up any leftover cluster role bindings
	cleanUpClusterRoleBindings(sandboxName)
	// Clean up any leftover cluster role bindings (do not delete sandbox pod bindings)
	cleanUpClusterRoleBindings(sandboxName, false)

	// ticker := time.NewTicker(retryTimerDuration * time.Millisecond)

@@ -394,8 +426,8 @@ func terminateScenario(sandboxName string, scenarioName string) {
func createSandbox(sandboxName string) {
	var err error

	// Clean up any leftover cluster role bindings first
	cleanUpClusterRoleBindings(sandboxName)
	// Clean up any leftover cluster role bindings first (clean all including old sandbox pod bindings)
	cleanUpClusterRoleBindings(sandboxName, true)

	// Create new Model instance
	modelCfg := mod.ModelCfg{
@@ -426,8 +458,8 @@ func destroySandbox(sandboxName string) {
	ve.activeScenarioNames[sandboxName] = ""
	ve.activeModels[sandboxName] = nil

	// Clean up any leftover cluster role bindings
	cleanUpClusterRoleBindings(sandboxName)
	// Clean up any leftover cluster role bindings (clean all when destroying sandbox)
	cleanUpClusterRoleBindings(sandboxName, true)

	// ticker := time.NewTicker(retryTimerDuration * time.Millisecond)

@@ -491,17 +523,17 @@ func deleteReleases(sandboxName string, scenarioName string, procName string) (e
			}
		}

		// Then delete charts
		// Then delete charts (queued sequentially in Helm worker)
		if _, err := os.Stat(path); err == nil {
			log.Debug("Removing charts from path: ", path)
			os.RemoveAll(path)
			log.Debug("Queueing chart removal from path: ", path)
			_ = helm.CleanChartDir(path)
		}
	}
	return err, chartsToDelete
}

func cleanUpClusterRoleBindings(sandboxName string) {
	log.Info("Cleaning up ClusterRoleBindings for sandbox: ", sandboxName)
func cleanUpClusterRoleBindings(sandboxName string, cleanAll bool) {
	log.Info("Cleaning up ClusterRoleBindings for sandbox: ", sandboxName, " (cleanAll=", cleanAll, ")")
	cmd := exec.Command("kubectl", "get", "clusterrolebindings", "-o", "name")
	out, err := cmd.Output()
	if err != nil {
@@ -509,11 +541,23 @@ func cleanUpClusterRoleBindings(sandboxName string) {
		return
	}

	sboxPods := strings.Split(strings.TrimSpace(os.Getenv("MEEP_SANDBOX_PODS")), ",")
	sboxPodMap := make(map[string]bool)
	for _, pod := range sboxPods {
		sboxPodMap[strings.TrimSpace(pod)] = true
	}

	lines := strings.Split(string(out), "\n")
	prefix := "clusterrolebinding.rbac.authorization.k8s.io/" + sandboxName + ":"
	for _, line := range lines {
		line = strings.TrimSpace(line)
		if strings.HasPrefix(line, prefix) {
			if !cleanAll {
				podName := strings.TrimPrefix(line, prefix)
				if sboxPodMap[podName] {
					continue
				}
			}
			log.Info("Deleting leftover clusterrolebinding: ", line)
			deleteCmd := exec.Command("kubectl", "delete", line)
			_ = deleteCmd.Run()
Loading