Fixed new microservice-worker architecture to not duplicate executions

This commit is contained in:
frikky
2021-12-21 00:17:13 +01:00
parent 9e8cc0e91e
commit bd248ef7f1
4 changed files with 54 additions and 39 deletions
+2 -1
View File
@@ -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)
}
}()
+47 -37
View File
@@ -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}`)))
}