From 35cf1132c8962702833729989c3c0db46283bd5a Mon Sep 17 00:00:00 2001 From: Frikky Date: Fri, 31 May 2024 14:57:27 +0200 Subject: [PATCH] Updated orborus & worker to use config the same way --- functions/onprem/orborus/go.mod | 4 +- functions/onprem/orborus/orborus.go | 2 +- functions/onprem/worker/go.mod | 2 +- functions/onprem/worker/worker.go | 165 +++++++++++----------------- 4 files changed, 66 insertions(+), 107 deletions(-) diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod index abf409b2..9101801f 100644 --- a/functions/onprem/orborus/go.mod +++ b/functions/onprem/orborus/go.mod @@ -4,13 +4,13 @@ go 1.22.0 toolchain go1.22.2 -replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared +//replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared require ( github.com/docker/docker v26.1.0+incompatible github.com/docker/go-connections v0.5.0 github.com/satori/go.uuid v1.2.0 - github.com/shuffle/shuffle-shared v0.6.29 + github.com/shuffle/shuffle-shared v0.6.37 k8s.io/api v0.30.0 k8s.io/apimachinery v0.30.0 k8s.io/client-go v0.30.0 diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 4b291a6b..60a82a1f 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -678,7 +678,7 @@ func deployWorker(image string, identifier string, env []string, executionReques } - clientset, config, err := shuffle.GetKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Printf("[ERROR] Error getting kubernetes client:", err) return err diff --git a/functions/onprem/worker/go.mod b/functions/onprem/worker/go.mod index 248a40f5..c6bfe36b 100644 --- a/functions/onprem/worker/go.mod +++ b/functions/onprem/worker/go.mod @@ -6,7 +6,7 @@ require ( github.com/docker/docker v26.1.0+incompatible github.com/gorilla/mux v1.8.1 github.com/satori/go.uuid v1.2.0 - github.com/shuffle/shuffle-shared v0.6.30 + github.com/shuffle/shuffle-shared v0.6.37 k8s.io/api v0.30.0 k8s.io/apimachinery v0.30.0 k8s.io/client-go v0.30.0 diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index fe145066..7ecbed96 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -3,6 +3,7 @@ package main import ( "github.com/shuffle/shuffle-shared" + "bytes" "context" "encoding/json" @@ -21,8 +22,8 @@ import ( "time" "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/container" "github.com/docker/docker/api/types/mount" dockerclient "github.com/docker/docker/client" // This is for automatic removal of certain code :) @@ -34,10 +35,6 @@ import ( corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" - "k8s.io/client-go/rest" - "k8s.io/client-go/tools/clientcmd" - "k8s.io/client-go/util/homedir" - "path/filepath" ) // This is getting out of hand :) @@ -56,7 +53,6 @@ var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") // var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME") - // var baseimagename = "registry.hub.docker.com/frikky/shuffle" var registryName = "registry.hub.docker.com" var sleepTime = 2 @@ -81,7 +77,6 @@ var startAction string //var allLogs map[string]string //var containerIds []string var downloadedImages []string - type ImageDownloadBody struct { Image string `json:"image"` } @@ -93,6 +88,7 @@ type ImageRequest struct { var finishedExecutions []string var imagesDistributed []string + // Images to be autodeployed in the latest version of Shuffle. var autoDeploy = map[string]string{ "http:1.4.0": "frikky/shuffle:http_1.4.0", @@ -139,6 +135,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo return err } + handleExecutionResult(workflowExecution) validated := shuffle.ValidateFinished(ctx, -1, workflowExecution) if validated { @@ -178,7 +175,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo } } - if len(subflowId) == 0 { + if len(subflowId) == 0 { log.Printf("[DEBUG][%s] No waiting result found. Not polling", workflowExecution.ExecutionId) for _, action := range workflowExecution.Workflow.Actions { @@ -186,17 +183,19 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo workflowExecution.Workflow.Triggers = append(workflowExecution.Workflow.Triggers, shuffle.Trigger{ AppName: action.AppName, Parameters: action.Parameters, - ID: action.ID, + ID: action.ID, }) } } + for _, trigger := range workflowExecution.Workflow.Triggers { //log.Printf("[DEBUG] Found trigger %s", trigger.AppName) if trigger.AppName != "User Input" && trigger.AppName != "Shuffle Workflow" && trigger.AppName != "shuffle-subflow" { continue } + // check if it has wait for results in params wait := false for _, param := range trigger.Parameters { @@ -215,9 +214,9 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo //log.Printf("[DEBUG][%s] Found result %s", workflowExecution.ExecutionId, result.Action.ID) if result.Action.ID == trigger.ID && result.Status != "SUCCESS" && result.Status != "FAILURE" { //log.Printf("[DEBUG][%s] Found subflow result that is not handled. Waiting for results", workflowExecution.ExecutionId) - + subflowId = result.Action.ID - found = true + found = true break } } @@ -236,20 +235,21 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo if len(subflowId) > 0 { // Under rerun period timeout - timeComparison := 120 + timeComparison := 120 log.Printf("[DEBUG][%s] Starting polling for %d seconds to see if new subflow updates are found on the backend that are not handled. Subflow ID: %s", workflowExecution.ExecutionId, timeComparison, subflowId) timestart := time.Now() streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl) for { - err = handleSubflowPoller(ctx, workflowExecution, streamResultUrl, subflowId) + err = handleSubflowPoller(ctx, workflowExecution, streamResultUrl, subflowId) if err == nil { log.Printf("[DEBUG] Subflow is finished and we are breaking the thingy") - + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" && workflowExecution.ExecutionSource != "default" { log.Printf("[DEBUG] Force shutdown of worker due to optimized run with webserver. Expecting reruns to take care of this") os.Exit(0) } + break } @@ -271,6 +271,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo return nil } + // removes every container except itself (worker) func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason string, handleResultSend bool) { 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) @@ -310,7 +311,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason } */ } else { - + } if len(reason) > 0 && len(nodeId) > 0 { @@ -393,7 +394,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } } - clientset, err := getKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Printf("[ERROR] Failed getting kubernetes: %s", err) return err @@ -499,7 +500,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] if !strings.Contains(param.Value, "shuffle-backend") { continue - } + } // Automatic replacement as this is default if len(os.Getenv("BASE_URL")) > 0 { @@ -514,6 +515,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } } + // Max 10% CPU every second //CPUShares: 128, //CPUQuota: 10000, @@ -540,7 +542,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] // Get environment for certificates volumeBinds := []string{} - volumeBindString := os.Getenv("SHUFFLE_VOLUME_BINDS") + volumeBindString:= os.Getenv("SHUFFLE_VOLUME_BINDS") if len(volumeBindString) > 0 { volumeBindSplit := strings.Split(volumeBindString, ",") for _, volumeBind := range volumeBindSplit { @@ -581,6 +583,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] Env: env, } + // Checking as late as possible, just in case. newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID) _, err := shuffle.GetCache(ctx, newExecId) @@ -860,7 +863,7 @@ func askOtherWorkersToDownloadImage(image string) { // Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" { log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. NOT distributing images %s", image) - return + return } if shuffle.ArrayContains(imagesDistributed, image) { @@ -893,7 +896,7 @@ func askOtherWorkersToDownloadImage(image string) { req, err := http.NewRequest( "POST", url, - bytes.NewBuffer(imageJSON), + bytes.NewBuffer(imageJSON), ) if err != nil { @@ -933,6 +936,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { return } + startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId) dockercli, err := dockerclient.NewEnvClient() @@ -1000,7 +1004,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { // marshal action and put it in there rofl //log.Printf("[INFO][%s] Time to execute %s (%s) with app %s:%s, function %s, env %s with %d parameters.", workflowExecution.ExecutionId, action.ID, action.Label, action.AppName, action.AppVersion, action.Name, action.Environment, len(action.Parameters)) - + log.Printf("[DEBUG][%s] Action: Send, Label: '%s', Action: '%s', Run status: %s, Extra=", workflowExecution.ExecutionId, action.Label, action.AppName, workflowExecution.Status) actionData, err := json.Marshal(action) @@ -1086,9 +1090,10 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { } if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 { - env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT"))) + env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT"))) } + // Fixes issue: // standard_go init_linux.go:185: exec user process caused "argument list too long" // https://devblogs.microsoft.com/oldnewthing/20100203-00/?p=15083 @@ -1116,6 +1121,8 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { fmt.Sprintf("%s:%s_%s", baseimagename, parsedAppname, action.AppVersion), } + + // If cleanup is set, it should run for efficiency pullOptions := types.ImagePullOptions{} if cleanupEnv == "true" { @@ -1378,7 +1385,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { log.Printf("[DEBUG][%s] Shutting down (17)", workflowExecution.ExecutionId) if isKubernetes == "true" { // log.Printf("workflow execution: %#v", workflowExecution) - clientset, err := getKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Println("[ERROR] Error getting kubernetes client (1):", err) os.Exit(1) @@ -1587,7 +1594,7 @@ func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.Workflow log.Printf("[DEBUG] Shutting down (20)") if isKubernetes == "true" { // log.Printf("workflow execution: %#v", workflowExecution) - clientset, err := getKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Println("[ERROR] Error getting kubernetes client (2):", err) os.Exit(1) @@ -1618,6 +1625,7 @@ func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.Workflow } } + if workflowExecution.Status == "WAITING" && workflowExecution.ExecutionSource != "default" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" { log.Printf("[INFO][%s] Workflow execution is waiting. Exiting worker, as backend will restart it.", workflowExecution.ExecutionId) shutdown(workflowExecution, "", "", true) @@ -1683,7 +1691,7 @@ func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffl log.Printf("[DEBUG] Shutting down (20)") if isKubernetes == "true" { // log.Printf("workflow execution: %#v", workflowExecution) - clientset, err := getKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Println("[ERROR] Error getting kubernetes client (2):", err) os.Exit(1) @@ -1700,7 +1708,7 @@ func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffl log.Printf("[DEBUG] Shutting down (21)") if isKubernetes == "true" { // log.Printf("workflow execution: %#v", workflowExecution) - clientset, err := getKubernetesClient() + clientset, _, err := shuffle.GetKubernetesClient() if err != nil { log.Println("[ERROR] Error getting kubernetes client (3):", err) os.Exit(1) @@ -1888,62 +1896,6 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { return envVars } -func getKubernetesClient() (*kubernetes.Clientset, error) { - - // Gets the config content from Orborus. - kubeconfigContent := os.Getenv("KUBERNETES_CONFIG") - if len(kubeconfigContent) > 0 { - log.Printf("[INFO] Using KUBERNETES_CONFIG to set up Kubernetes client: %#v", os.Getenv("KUBERNETES_CONFIG")) - config, err := rest.InClusterConfig() - if err != nil { - log.Printf("[ERROR] Failed to create Kubernetes client from in-cluster config: %s", err) - } else { - // Replace client configuration with kubeconfig content - config, err = clientcmd.RESTConfigFromKubeConfig([]byte(kubeconfigContent)) - if err != nil { - log.Printf("[ERROR] Failed to create Kubernetes client from KUBERNETES_CONFIG: %s", err) - } else { - // Create Kubernetes client - clientset, err := kubernetes.NewForConfig(config) - if err != nil { - return nil, err - } - - return clientset, nil - } - } - } - - // Fallback - if isRunningInCluster() { - config, err := rest.InClusterConfig() - if err != nil { - return nil, err - } - - clientset, err := kubernetes.NewForConfig(config) - if err != nil { - return nil, err - } - - return clientset, nil - } - - home := homedir.HomeDir() - kubeconfigPath := filepath.Join(home, ".kube", "config") - config, err := clientcmd.BuildConfigFromFlags("", kubeconfigPath) - if err != nil { - return nil, err - } - - clientset, err := kubernetes.NewForConfig(config) - if err != nil { - return nil, err - } - - return clientset, nil -} - func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { if request.Body == nil { resp.WriteHeader(http.StatusBadRequest) @@ -2050,9 +2002,10 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl resp.Write([]byte(fmt.Sprintf(`{"success": true, "reason": "Execution is not executing, but %s"}`, workflowExecution.Status))) } + log.Printf("[DEBUG][%s] Shutting down (35)", workflowExecution.ExecutionId) - // Force sending result + // Force sending result shutdownData, err := json.Marshal(workflowExecution) if err != nil { log.Printf("[ERROR][%s] Failed marshalling execution (35): %s", workflowExecution.ExecutionId, err) @@ -2128,11 +2081,11 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl attempts += 1 log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(parsedValue.Results), resultLength) /* - if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) { - log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)) - runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp) - return - } + if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) { + log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)) + runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp) + return + } */ } } @@ -2165,6 +2118,8 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } func sendSelfRequest(actionResult shuffle.ActionResult) { + + data, err := json.Marshal(actionResult) if err != nil { log.Printf("[ERROR][%s] Shutting down (24): Failed to unmarshal data for backend: %s", actionResult.ExecutionId, err) @@ -2223,10 +2178,10 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { // Basically to reduce backend strain /* - if shuffle.ArrayContains(finishedExecutions, workflowExecution.ExecutionId) { - log.Printf("[INFO][%s] NOT sending backend info since it's already been sent before.", workflowExecution.ExecutionId) - return - } + if shuffle.ArrayContains(finishedExecutions, workflowExecution.ExecutionId) { + log.Printf("[INFO][%s] NOT sending backend info since it's already been sent before.", workflowExecution.ExecutionId) + return + } */ // Take it down again @@ -2237,7 +2192,6 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { } finishedExecutions = append(finishedExecutions, workflowExecution.ExecutionId) - */ streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl) @@ -2294,7 +2248,6 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || (len(environments) == 1 && requestsSent == 0 && len(workflowExecution.Results) >= 1 && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm") || (len(workflowExecution.Results) >= len(workflowExecution.Workflow.Actions)+extra && len(workflowExecution.Workflow.Actions) > 0) { - if workflowExecution.Status == "FINISHED" { for _, result := range workflowExecution.Results { if result.Status == "EXECUTING" || result.Status == "WAITING" { @@ -2303,7 +2256,8 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { } } } - + + log.Printf("[DEBUG][%s] Should send full result to %s", workflowExecution.ExecutionId, baseUrl) //data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization) @@ -2387,7 +2341,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { // GetLocalIP returns the non loopback local IP of the host func getLocalIP() string { - + addrs, err := net.InterfaceAddrs() if err != nil { return "" @@ -2432,7 +2386,8 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { } log.Printf("[DEBUG] OLD HOSTNAME: %s", appCallbackUrl) - + + port := listener.Addr().(*net.TCPAddr).Port // Set the port environment variable os.Setenv("WORKER_PORT", fmt.Sprintf("%d", port)) @@ -2447,22 +2402,21 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { func downloadDockerImageBackend(client *http.Client, imageName string) error { // Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" { - //log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. Not downloading image %s", imageName) + log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. Not downloading image %s", imageName) return nil } if arrayContains(downloadedImages, imageName) { - log.Printf("[DEBUG] Image %s already downloaded", imageName) + log.Printf("[DEBUG] Image %s already downloaded - not re-downloading", imageName) return nil } + log.Printf("[DEBUG] Trying to download image %s from backend %s as it doesn't exist. All images: %#v", imageName, baseUrl, downloadedImages) downloadedImages = append(downloadedImages, imageName) data := fmt.Sprintf(`{"name": "%s"}`, imageName) dockerImgUrl := fmt.Sprintf("%s/api/v1/get_docker_image", baseUrl) - - log.Printf("[DEBUG] Trying to download image %s from backend %s as it doesn't exist. Data sent: %#v, All images: %#v", imageName, baseUrl, data, downloadedImages) req, err := http.NewRequest( "POST", @@ -2591,6 +2545,7 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { */ } + // Runs data discovery func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, action *shuffle.Action, workflowExecution *shuffle.WorkflowExecution) error { @@ -2854,7 +2809,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe if newresp.StatusCode != 200 { log.Printf("[ERROR] %sStatusCode (1): %d", string(body), newresp.StatusCode) time.Sleep(time.Duration(sleepTime) * time.Second) - return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode)) + return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode) ) } err = json.Unmarshal(body, &workflowExecution) @@ -2937,6 +2892,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe // Set environment variable + //log.Printf("Before wait") //wg := sync.WaitGroup{} //wg.Add(1) @@ -3008,6 +2964,7 @@ func main() { swarmConfig := os.Getenv("SHUFFLE_SWARM_CONFIG") log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, swarmConfig) + authorization := "" executionId := "" @@ -3312,6 +3269,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) { return } + for _, img := range images { for _, tag := range img.RepoTags { splitTag := strings.Split(tag, ":") @@ -3324,7 +3282,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) { possibleNames = append(possibleNames, fmt.Sprintf("frikky/shuffle:%s", baseTag)) possibleNames = append(possibleNames, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", baseTag)) - if arrayContains(possibleNames, image.Image) { + if (arrayContains(possibleNames, image.Image)) { log.Printf("[DEBUG] Image %s already downloaded that has been requested to download", image.Image) resp.WriteHeader(200) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "image already present"}`))) @@ -3349,6 +3307,7 @@ func runWebserver(listener net.Listener) { r.HandleFunc("/api/v1/run", handleRunExecution).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/download", handleDownloadImage).Methods("POST", "OPTIONS") + if strings.ToLower(os.Getenv("SHUFFLE_DEBUG_MEMORY")) == "true" { r.HandleFunc("/debug/pprof/", pprof.Index) r.HandleFunc("/debug/pprof/heap", pprof.Handler("heap").ServeHTTP)