diff --git a/functions/k8s-confs/env-configmap.yaml b/functions/k8s-confs/env-configmap.yaml index fb8d9d28..ad02ad97 100644 --- a/functions/k8s-confs/env-configmap.yaml +++ b/functions/k8s-confs/env-configmap.yaml @@ -53,7 +53,7 @@ data: SSO_REDIRECT_URL: "" TZ: "Europe/Amsterdam \t\t\t\t\t" IS_KUBERNETES: "true" - REGISTRY_URL: "shuffle-registry:5000" + REGISTRY_URL: "172.17.14.71:5000" #REGISTRY_AUTH: "" kind: ConfigMap metadata: diff --git a/functions/onprem/worker/Dockerfile b/functions/onprem/worker/Dockerfile index 8d3f6155..b1bbbe43 100755 --- a/functions/onprem/worker/Dockerfile +++ b/functions/onprem/worker/Dockerfile @@ -32,6 +32,8 @@ FROM alpine:3.15.0 ENV SHUFFLE_BASE_IMAGE_REGISTRY=docker.io ENV SHUFFLE_BASE_IMAGE_NAME=frikky/shuffle ENV SHUFFLE_BASE_IMAGE_TAG_SUFFIX=0.8.70 + +#for k8s ENV SHUFFLE_OPENSEARCH_URL=https://opensearch:9200 ENV SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY=true diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index fe47bfac..cc7dcbc4 100755 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -22,7 +22,7 @@ import ( "github.com/docker/docker/api/types" "github.com/docker/docker/api/types/container" //"github.com/docker/docker/api/types/filters" - // "github.com/docker/docker/api/types/mount" + "github.com/docker/docker/api/types/mount" dockerclient "github.com/docker/docker/client" //"github.com/go-git/go-billy/v5/memfs" @@ -112,18 +112,6 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason log.Printf("[DEBUG][%s] Shutdown (%s) started with reason %#v. Result amount: %d. ResultsSent: %d, Send result: %#v, Parent: %#v", workflowExecution.ExecutionId, workflowExecution.Status, reason, len(workflowExecution.Results), requestsSent, handleResultSend, workflowExecution.ExecutionParent) //reason := "Error in execution" - // if os.Getenv("IS_KUBERNETES") == "true" { - // log.Printf("[DEBUG][%s] Running in Kubernetes shutting down") - // workerPod := fmt.Sprintf("worker-%s", workflowExecution.ExecutionId) - // namespace := "shuffle" - // clientset, err := getKubernetesClient() - // if err != nil { - // log.Printf("[ERROR] Error getting kubernetes client: %s", err) - // return - // } - - // } else { - sleepDuration := 1 if handleResultSend && requestsSent < 2 { shutdownData, err := json.Marshal(workflowExecution) @@ -237,12 +225,10 @@ func getKubernetesClient() (*kubernetes.Clientset, error) { // Deploys the internal worker whenever something happens func deployApp(cli *dockerclient.Client, image string, identifier string, env []string, workflowExecution shuffle.WorkflowExecution, action shuffle.Action) error { - log.Printf("HELLO FROM DEPLOYAPP") + log.Printf("################################### new call to deployApp ###################################") log.Printf("image: %s", image) log.Printf("identifier: %s", identifier) - - log.Printf("IS_KUBERNETES: %s", os.Getenv("IS_KUBERNETES")) - log.Printf("REGISTRY_NAME: %s", os.Getenv("REGISTRY_NAME")) + // log.Printf("execution: %+v", workflowExecution) if os.Getenv("IS_KUBERNETES") == "true" { @@ -268,6 +254,26 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] value := strSplit[0] value = strings.ReplaceAll(value, "_", "-") + // checking if app is generated or not + appDetails := strings.Split(image, ":")[1] + appDetailsSplit := strings.Split(appDetails, "_") + appName := appDetailsSplit[0] + appVersion := appDetailsSplit[1] + + for _, app := range workflowExecution.Workflow.Actions { + log.Printf("[DEBUG] App: %s, Version: %s", appName, appVersion) + log.Printf("[DEBUG] Checking app %s with version %s", app.AppName, app.AppVersion) + if app.AppName == appName && app.AppVersion == appVersion { + if app.Generated == true { + log.Printf("[DEBUG] Generated app, setting local registry") + image = fmt.Sprintf("%s/%s", registryName, image) + break + } else { + log.Printf("[DEBUG] Not generated app, setting shuffle registry") + } + } + } + //fix naming convention podUuid := uuid.NewV4().String() podName := fmt.Sprintf("%s-%s", value, podUuid) @@ -297,11 +303,112 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] fmt.Fprintf(os.Stderr, "Error creating pod: %v\n", err) os.Exit(1) } - fmt.Printf("Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace) + fmt.Printf("[DEBUG] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace) } else { - // docker part + // form basic hostConfig + ctx := context.Background() + + if action.AppName == "shuffle-subflow" { + // Automatic replacement of URL + for paramIndex, param := range action.Parameters { + if param.Name != "backend_url" { + continue + } + + if strings.Contains(param.Value, "shuffle-backend") { + // Automatic replacement as this is default + action.Parameters[paramIndex].Value = os.Getenv("BASE_URL") + log.Printf("[DEBUG][%s] Replaced backend_url with %s", workflowExecution.ExecutionId, os.Getenv("BASE_URL")) + } + } } + // Max 10% CPU every second + //CPUShares: 128, + //CPUQuota: 10000, + //CPUPeriod: 100000, + hostConfig := &container.HostConfig{ + LogConfig: container.LogConfig{ + Type: "json-file", + Config: map[string]string{ + "max-size": "10m", + }, + }, + Resources: container.Resources{}, + } + + hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId)) + + // Removing because log extraction should happen first + if cleanupEnv == "true" { + hostConfig.AutoRemove = true + } + + // FIXME: Add proper foldermounts here + //log.Printf("\n\nPRE FOLDERMOUNT\n\n") + //volumeBinds := []string{"/tmp/shuffle-mount:/rules"} + //volumeBinds := []string{"/tmp/shuffle-mount:/rules"} + volumeBinds := []string{} + if len(volumeBinds) > 0 { + log.Printf("[DEBUG] Setting up binds for container!") + hostConfig.Binds = volumeBinds + hostConfig.Mounts = []mount.Mount{} + for _, bind := range volumeBinds { + if !strings.Contains(bind, ":") || strings.Contains(bind, "..") || strings.HasPrefix(bind, "~") { + log.Printf("[WARNING] Bind %s is invalid.", bind) + continue + } + + log.Printf("[DEBUG] Appending bind %s", bind) + bindSplit := strings.Split(bind, ":") + sourceFolder := bindSplit[0] + destinationFolder := bindSplit[0] + hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{ + Type: mount.TypeBind, + Source: sourceFolder, + Target: destinationFolder, + }) + } + } else { + //log.Printf("[WARNING] Not mounting folders") + } + + config := &container.Config{ + Image: image, + Env: env, + } + + // Checking as late as possible, just in case. + newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID) + _, err := shuffle.GetCache(ctx, newExecId) + if err == nil { + log.Printf("\n\n[DEBUG] Result for %s already found - returning\n\n", newExecId) + return nil + } + + cacheData := []byte("1") + err = shuffle.SetCache(ctx, newExecId, cacheData, 30) + if err != nil { + log.Printf("[WARNING] Failed setting cache for action %s: %s", newExecId, err) + } else { + log.Printf("[DEBUG] Adding %s to cache. Name: %s", newExecId, action.Name) + } + + if action.ExecutionDelay > 0 { + log.Printf("[DEBUG] Running app %s in docker with delay of %d", action.Name, action.ExecutionDelay) + waitTime := time.Duration(action.ExecutionDelay) * time.Second + + time.AfterFunc(waitTime, func() { + DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId) + }) + } else { + log.Printf("[DEBUG] Running app %s in docker NORMALLY as there is no delay set with identifier %s", action.Name, identifier) + returnvalue := DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId) + log.Printf("[DEBUG] Normal deploy ret: %s", returnvalue) + return returnvalue + } + return nil + } return nil } @@ -1168,7 +1275,17 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec if workflowExecution.Status != "EXECUTING" { log.Printf("[WARNING][%s] Exiting as worker execution has status %s!", workflowExecution.ExecutionId, workflowExecution.Status) log.Printf("[DEBUG] Shutting down (21)") - shutdown(workflowExecution, "", "", true) + if os.Getenv("IS_KUBERNETES") == "true" { + // log.Printf("workflow execution: %#v", workflowExecution) + clientset, err := getKubernetesClient() + if err != nil { + fmt.Println("[ERROR]Error getting kubernetes client:", err) + os.Exit(1) + } + cleanupExecution(clientset, workflowExecution, "shuffle") + } else { + shutdown(workflowExecution, "", "", true) + } } setWorkflowExecution(ctx, workflowExecution, false)