From bd248ef7f1aa08be2679e0986e3ccb606a7c2fd5 Mon Sep 17 00:00:00 2001 From: frikky Date: Tue, 21 Dec 2021 00:17:13 +0100 Subject: [PATCH] Fixed new microservice-worker architecture to not duplicate executions --- backend/app_sdk/app_base.py | 4 ++ backend/go-app/walkoff.go | 2 +- functions/onprem/orborus/orborus.go | 3 +- functions/onprem/worker/worker.go | 84 ++++++++++++++++------------- 4 files changed, 54 insertions(+), 39 deletions(-) diff --git a/backend/app_sdk/app_base.py b/backend/app_sdk/app_base.py index c205c722..5df2d869 100644 --- a/backend/app_sdk/app_base.py +++ b/backend/app_sdk/app_base.py @@ -78,6 +78,10 @@ class AppBase: self.logger.info("[DEBUG] Value too large. Returning from magic") return input_data + if not "\n" in input_data and not "," in input_data: + self.logger.info("[DEBUG] No data to autoparse") + return input_data + new_input = input_data try: #new_input.strip() diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 9d45f0cc..ddbb3ce7 100644 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -412,7 +412,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { } //log.Printf("Actionresult unmarshal: %s", string(body)) - log.Printf("[DEBUG] Got workflow result from %s of length %d", request.RemoteAddr, len(body)) + log.Printf("[DEBUG] Got workflow result from %s of length %d. \n\nBody: %s\n\n", request.RemoteAddr, len(body), string(body)) err = shuffle.ValidateNewWorkerExecution(body) if err == nil { resp.WriteHeader(200) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index dcf97efe..6f9b2d85 100644 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -393,7 +393,8 @@ func deployWorker(image string, identifier string, env []string, executionReques } if err == nil { - executionIds = append(executionIds, executionRequest.ExecutionId) + // FIXME: Readd this? Removed for rerun reasons + // executionIds = append(executionIds, executionRequest.ExecutionId) } }() diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index bb78b418..51c95621 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -137,8 +137,9 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason } */ } else { - log.Printf("[DEBUG][%s] NOT cleaning up containers. IDS: %d, CLEANUP env: %s", workflowExecution.ExecutionId, len(containerIds), cleanupEnv) - + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" { + log.Printf("[DEBUG][%s] NOT cleaning up containers. IDS: %d, CLEANUP env: %s", workflowExecution.ExecutionId, len(containerIds), cleanupEnv) + } } if len(reason) > 0 && len(nodeId) > 0 { @@ -205,7 +206,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason log.Printf("[WARNING][%s] Failed abort request: %s", workflowExecution.ExecutionId, err) } } else { - log.Printf("[INFO][%s] NOT running abort during shutdown.", workflowExecution.ExecutionId) + //log.Printf("[INFO][%s] NOT running abort during shutdown.", workflowExecution.ExecutionId) } log.Printf("[INFO][%s] Finished shutdown (after %d seconds). ", workflowExecution.ExecutionId, sleepDuration) @@ -216,7 +217,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason time.Sleep(time.Duration(sleepDuration) * time.Second) os.Exit(3) } else { - log.Printf("\n\n[DEBUG][%s] Sending result and resetting values (K8s & Swarm).\n\n", workflowExecution.ExecutionId) + log.Printf("[DEBUG][%s] Sending result and resetting values (K8s & Swarm).", workflowExecution.ExecutionId) //UpdateExecutionVariables(ctx, workflowExecution.ExecutionId, startAction, children, parents, visited, executed, nextActions, environments, extra) /* @@ -232,7 +233,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason results = []shuffle.ActionResult{} allLogs = map[string]string{} */ - requestsSent = 0 + //requestsSent = 0 //executionRunning = false } //cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId) @@ -249,7 +250,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] appName := strings.Replace(identifier, fmt.Sprintf("_%s", action.ID), "", -1) appName = strings.Replace(appName, fmt.Sprintf("_%s", workflowExecution.ExecutionId), "", -1) appName = strings.ToLower(appName) - log.Printf("[INFO][%s] New appname: %s, image: %s", workflowExecution.ExecutionId, appName, image) + //log.Printf("[INFO][%s] New appname: %s, image: %s", workflowExecution.ExecutionId, appName, image) if !shuffle.ArrayContains(downloadedImages, image) { log.Printf("[DEBUG] Downloading image %s from backend as it's first iteration for this image on the worker.", image) @@ -270,14 +271,14 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] return err } - log.Printf("[DEBUG][%s] Should run towards port %d for app %s", workflowExecution.ExecutionId, exposedPort, appName) + //log.Printf("[DEBUG][%s] Should run towards port %d for app %s", workflowExecution.ExecutionId, exposedPort, appName) err = sendAppRequest(baseUrl, appName, exposedPort, action, workflowExecution) if err != nil { log.Printf("[ERROR] Failed sending request to app %s on port %d: %s", appName, exposedPort, err) return err } - log.Printf("[DEBUG] Successfully ran request towards port %d for app %s", exposedPort, appName) + //log.Printf("[DEBUG] Successfully ran request towards port %d for app %s", exposedPort, appName) return nil } @@ -763,7 +764,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { } if exit && len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions) { - log.Printf("[DEBUG] Shutting down (1)") + log.Printf("[DEBUG][%s] Shutting down (1)", workflowExecution.ExecutionId) shutdown(workflowExecution, "", "", true) } @@ -1105,7 +1106,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { // Sending full execution so that it won't have to load in every app // This might be an issue if they can read environments, but that's alright // if everything is generated during execution - log.Printf("[INFO][%s] Deployed with CALLBACK_URL %s and BASE_URL %s", workflowExecution.ExecutionId, appCallbackUrl, baseUrl) + //log.Printf("[DEBUG][%s] Deployed with CALLBACK_URL %s and BASE_URL %s", workflowExecution.ExecutionId, appCallbackUrl, baseUrl) env := []string{ fmt.Sprintf("ACTION=%s", string(actionData)), fmt.Sprintf("EXECUTIONID=%s", workflowExecution.ExecutionId), @@ -1901,7 +1902,6 @@ func runTestExecution(client *http.Client, workflowId, apikey string) (string, s } func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { - log.Printf("\n\n[DEBUG] In workflowQueue\n\n") body, err := ioutil.ReadAll(request.Body) if err != nil { log.Printf("[WARNING] (3) Failed reading body for workflowqueue") @@ -1909,6 +1909,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return } + log.Printf("[DEBUG] In workflowQueue with body length %d", len(body)) //log.Printf("Got result: %s", string(body)) var actionResult shuffle.ActionResult @@ -1929,7 +1930,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { ctx := context.Background() workflowExecution, err := getWorkflowExecution(ctx, actionResult.ExecutionId) if err != nil { - log.Printf("[ERROR] Failed getting execution (workflowqueue) %s: %s", actionResult.ExecutionId, err) + log.Printf("[ERROR][%s] Failed getting execution (workflowqueue) %s: %s", actionResult.ExecutionId, actionResult.ExecutionId, err) resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting execution ID %s because it doesn't exist locally."}`, actionResult.ExecutionId))) return @@ -1963,18 +1964,14 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { results = append(results, actionResult) - resp.WriteHeader(200) - resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) - - log.Printf("\n\n[DEBUG] In workflowQueue with transaction\n\n") + log.Printf("[DEBUG][%s] In workflowQueue with transaction", workflowExecution.ExecutionId) runWorkflowExecutionTransaction(ctx, 0, workflowExecution.ExecutionId, actionResult, resp) } // Will make sure transactions are always ran for an execution. This is recursive if it fails. Allowed to fail up to 5 times func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workflowExecutionId string, actionResult shuffle.ActionResult, resp http.ResponseWriter) { - //log.Printf("IN WORKFLOWEXECUTION SUB!") - // Should start a tx for the execution here + log.Printf("[DEBUG][%s] IN WORKFLOWEXECUTION SUB!", actionResult.ExecutionId) workflowExecution, err := getWorkflowExecution(ctx, workflowExecutionId) if err != nil { log.Printf("[ERROR] Failed getting execution cache: %s", err) @@ -2017,11 +2014,11 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl return } } - //log.Printf(`[INFO] Got result %s from %s`, actionResult.Status, actionResult.Action.ID) + log.Printf(`[DEBUG][%s] Got result %s from %s. Execution status: %s. Save: %#v`, actionResult.ExecutionId, actionResult.Status, actionResult.Action.ID, workflowExecution.Status, dbSave) //dbSave := false if len(results) != len(workflowExecution.Results) { - log.Printf("[DEBUG] There may have been an issue in transaction queue. Result lengths: %d vs %d. Should check which exists the base results, but not in entire execution, then append.", len(results), len(workflowExecution.Results)) + log.Printf("[DEBUG][%s] There may have been an issue in transaction queue. Result lengths: %d vs %d. Should check which exists the base results, but not in entire execution, then append.", workflowExecution.ExecutionId, len(results), len(workflowExecution.Results)) } // Validating that action results hasn't changed @@ -2044,14 +2041,19 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } if setExecution || workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || workflowExecution.Status == "FAILURE" { + log.Printf("[INFO][%s] Running setexec with status %s", workflowExecution.ExecutionId, workflowExecution.Status) err = setWorkflowExecution(ctx, *workflowExecution, dbSave) if err != nil { resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed setting workflowexecution actionresult: %s"}`, err))) return } + + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + validateFinished(*workflowExecution) + } } else { - log.Printf("[INFO] Skipping setexec with status %s", workflowExecution.Status) + log.Printf("[INFO][%s] Skipping setexec with status %s", workflowExecution.ExecutionId, workflowExecution.Status) // Just in case. Should MAYBE validate finishing another time as well. // This fixes issues with e.g. shuffle.Action -> shuffle.Trigger -> shuffle.Action. @@ -2060,11 +2062,12 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } //if newExecutions && len(nextActions) > 0 { - // handleExecutionResult(*workflowExecution) + // log.Printf("[DEBUG][%s] New execution: %#v. NextActions: %#v", newExecutions, nextActions) + // //handleExecutionResult(*workflowExecution) //} - //resp.WriteHeader(200) - //resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) + resp.WriteHeader(200) + resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) } func getWorkflowExecution(ctx context.Context, id string) (*shuffle.WorkflowExecution, error) { @@ -2121,18 +2124,20 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) { //startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId) _, extra, _, _, _, _, _, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId) - log.Printf("[INFO] VALIDATION. Status: %s, shuffle.Actions: %d, Extra: %d, Results: %d\n", workflowExecution.Status, len(workflowExecution.Workflow.Actions), extra, len(workflowExecution.Results)) + log.Printf("[INFO][%s] VALIDATION. Status: %s, shuffle.Actions: %d, Extra: %d, Results: %d\n", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Workflow.Actions), extra, len(workflowExecution.Results)) //if len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions)+extra { if (len(environments) == 1 && requestsSent == 0 && len(workflowExecution.Results) >= 1) || (len(workflowExecution.Results) >= len(workflowExecution.Workflow.Actions) && len(workflowExecution.Workflow.Actions) > 0) { - requestsSent += 1 - log.Printf("[DEBUG] Should send full result to %s", baseUrl) + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" { + requestsSent += 1 + } + + log.Printf("[DEBUG][%s] Should send full result to %s", workflowExecution.ExecutionId, baseUrl) //data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization) shutdownData, err := json.Marshal(workflowExecution) if err != nil { - log.Printf("[ERROR] Failed to unmarshal data for backend") - log.Printf("[DEBUG] Shutting down (24)") + log.Printf("[ERROR][%s] Shutting down (24): Failed to unmarshal data for backend: %s", workflowExecution.ExecutionId, err) shutdown(workflowExecution, "", "", true) } @@ -2141,7 +2146,6 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) { } func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { - //log.Printf("[DEBUG] Got stream result") body, err := ioutil.ReadAll(request.Body) if err != nil { log.Printf("Failed reading body for stream result queue") @@ -2149,6 +2153,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return } + log.Printf("[DEBUG] In get stream results with body length %d", len(body)) var actionResult shuffle.ActionResult err = json.Unmarshal(body, &actionResult) @@ -2197,6 +2202,10 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId) requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration) + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + return nil + } + handleExecutionResult(workflowExecution) validateFinished(workflowExecution) @@ -2204,7 +2213,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo // The worker may not be running the backend hmm if dbSave { if workflowExecution.ExecutionSource == "default" { - log.Printf("[DEBUG] Shutting down (25)") + log.Printf("[DEBUG][%s] Shutting down (25)", workflowExecution.ExecutionId) shutdown(workflowExecution, "", "", true) //log.Printf("[INFO] Not sending backend info since source is default") //return @@ -2638,7 +2647,7 @@ func findAppInfo(image, name string) (int, error) { //log.Printf("[DEBUG] Portmappings: %#v", portMappings) if exposedPort >= 0 { - log.Printf("[INFO] Found service %s on port %d - no need to deploy another", name, exposedPort) + //log.Printf("[INFO] Found service %s on port %d - no need to deploy another", name, exposedPort) } else { // Increment by 1 for highest port if highest <= baseport { @@ -2700,7 +2709,7 @@ func sendAppRequest(incomingUrl, appName string, port int, action shuffle.Action parsedRequest.Url = fmt.Sprintf("%s:%d", parsedBaseurl, baseport) } - log.Printf("[DEBUG][%s] Should add a baseurl for the app to get back to: %s", workflowExecution.ExecutionId, parsedRequest.Url) + //log.Printf("[DEBUG][%s] Should add a baseurl for the app to get back to: %s", workflowExecution.ExecutionId, parsedRequest.Url) } // FIXME: Swapping because this was confusing during dev @@ -2713,7 +2722,7 @@ func sendAppRequest(incomingUrl, appName string, port int, action shuffle.Action if len(hostname) > 0 { parsedRequest.BaseUrl = fmt.Sprintf("http://%s:%d", hostname, baseport) //parsedRequest.BaseUrl = fmt.Sprintf("http://shuffle-workers:%d", baseport) - log.Printf("[DEBUG][%s] Changing hostname to local hostname in Docker network for WORKER URL: %s", workflowExecution.ExecutionId, parsedRequest.BaseUrl) + //log.Printf("[DEBUG][%s] Changing hostname to local hostname in Docker network for WORKER URL: %s", workflowExecution.ExecutionId, parsedRequest.BaseUrl) } data, err := json.Marshal(parsedRequest) @@ -3112,6 +3121,7 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return } + log.Printf("[DEBUG] In run execution with body length %d", len(body)) var execRequest shuffle.OrborusExecutionRequest err = json.Unmarshal(body, &execRequest) @@ -3249,15 +3259,15 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) { err = executionInit(workflowExecution) if err != nil { - log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err) - log.Printf("[DEBUG] Shutting down (30)") + log.Printf("[INFO][%s] Shutting down (30) - Workflow setup failed: %s", workflowExecution.ExecutionId, workflowExecution.ExecutionId, err) resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Error in execution init: %s"}`, err))) return //shutdown(workflowExecution, "", "", true) } - go handleExecutionResult(workflowExecution) + //go handleExecutionResult(workflowExecution) + handleExecutionResult(workflowExecution) resp.WriteHeader(200) resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) }