From f0a299e69adb09b8dfcd04bd000297d60241aea7 Mon Sep 17 00:00:00 2001 From: frikky Date: Tue, 12 Jan 2021 17:59:14 +0100 Subject: [PATCH] Worker a lot optimized --- functions/onprem/worker/worker.go | 264 +++++++++++++++++++++--------- 1 file changed, 185 insertions(+), 79 deletions(-) diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 029823c6..d77eb362 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -35,7 +35,9 @@ var sleepTime = 2 var requestCache *cache.Cache var topClient *http.Client var data string +var requestsSent = 0 +var environments []string var parents map[string][]string var children map[string][]string var visited []string @@ -889,7 +891,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] return err } - log.Printf("[INFO] Container %s is created", cont.ID) + log.Printf("[INFO] Container %s is created for %s", cont.ID, identifier) return nil } @@ -1043,6 +1045,11 @@ func handleSubworkflowExecution(client *http.Client, workflowExecution WorkflowE } } +func removeIndex(s []string, i int) []string { + s[len(s)-1], s[i] = s[i], s[len(s)-1] + return s[:len(s)-1] +} + func handleExecutionResult(workflowExecution WorkflowExecution) { if len(startAction) == 0 { startAction = workflowExecution.Start @@ -1052,6 +1059,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } + //log.Printf("NEXTACTIONS: %s", nextActions) queueNodes := []string{} //if len(nextActions) == 0 { // nextActions = append(nextActions, startAction) @@ -1105,7 +1113,14 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } - nextActions = children[item.Action.ID] + //if len(nextActions) == 0 { + //nextActions = append(nextActions, children[item.Action.ID]...) + for _, child := range children[item.Action.ID] { + if !arrayContains(nextActions, child) && !arrayContains(visited, child) && !arrayContains(visited, child) { + nextActions = append(nextActions, child) + } + } + if len(appendActions) > 0 { log.Printf("APPENDED NODES: %#v", appendActions) nextActions = append(nextActions, appendActions...) @@ -1113,6 +1128,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } + //log.Printf("Nextactions: %s", nextActions) // This is a backup in case something goes wrong in this complex hellhole. // Max default execution time is 5 minutes for now anyway, which should take // care if it gets stuck in a loop. @@ -1127,6 +1143,11 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } + if len(environments) == 1 { + log.Printf("[INFO] Should send results to the backend because environments are %s", environments) + validateFinished(workflowExecution) + } + if exit && len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions) { log.Printf("Shutting down.") shutdown(workflowExecution.ExecutionId, workflowExecution.Workflow.ID) @@ -1186,6 +1207,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } + //log.Printf("Checking nextactions: %s", nextActions) for _, node := range nextActions { nodeChildren := children[node] for _, child := range nodeChildren { @@ -1194,16 +1216,21 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { } } } - //log.Printf("NEXT: %s", nextActions) - //log.Printf("queueNodes: %s", queueNodes) // IF NOT VISITED && IN toExecuteOnPrem // SKIP if it's not onprem - for _, nextAction := range nextActions { + toRemove := []int{} + for index, nextAction := range nextActions { action := getAction(workflowExecution, nextAction, environment) // check visited and onprem if arrayContains(visited, nextAction) { - log.Printf("ALREADY VISITIED (%s): %s", action.Label, nextAction) + //log.Printf("ALREADY VISITIED (%s): %s", action.Label, nextAction) + toRemove = append(toRemove, index) + //nextActions = removeIndex(nextActions, index) + + //validateFinished(workflowExecution) + _ = index + continue } @@ -1281,6 +1308,13 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { break } + } else { + //log.Printf("Handling action %#v", action) + } + + if len(toRemove) > 0 { + //toRemove = []int{} + //for index, nextAction := range nextActions { } // Not really sure how this edgecase happens. @@ -1358,6 +1392,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { //return err return } + stats, err := dockercli.ContainerInspect(context.Background(), identifier) if err != nil || stats.ContainerJSONBase.State.Status != "running" { // REMOVE @@ -1371,6 +1406,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { //log.Printf("WHAT TO DO HERE?: %s", err) } } else if stats.ContainerJSONBase.State.Status == "running" { + //log.Printf(" continue } @@ -1406,11 +1442,13 @@ func handleExecutionResult(workflowExecution 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("Deployed with CALLBACK_URL %s and BASE_URL %s", appCallbackUrl, baseUrl) env := []string{ fmt.Sprintf("ACTION=%s", string(actionData)), fmt.Sprintf("EXECUTIONID=%s", workflowExecution.ExecutionId), fmt.Sprintf("AUTHORIZATION=%s", workflowExecution.Authorization), - fmt.Sprintf("CALLBACK_URL=%s", appCallbackUrl), + fmt.Sprintf("CALLBACK_URL=%s", baseUrl), + fmt.Sprintf("BASE_URL=%s", appCallbackUrl), } // Fixes issue: @@ -1418,10 +1456,10 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { // https://devblogs.microsoft.com/oldnewthing/20100203-00/?p=15083 maxSize := 32700 - len(string(actionData)) - 2000 if len(executionData) < maxSize { - log.Printf("ADDING FULL_EXECUTION because size is smaller than %d", maxSize) + log.Printf("[INFO] ADDING FULL_EXECUTION because size is smaller than %d", maxSize) env = append(env, fmt.Sprintf("FULL_EXECUTION=%s", string(executionData))) } else { - log.Printf("Skipping FULL_EXECUTION because size is larger than %d", maxSize) + log.Printf("[WARNING] Skipping FULL_EXECUTION because size is larger than %d", maxSize) } // Uses a few ways of getting / checking if an app is available @@ -1504,6 +1542,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { log.Printf("Unable to create docker client: %s", err) shutdown(workflowExecution.ExecutionId, workflowExecution.Workflow.ID) } + if len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions)+extra { shutdownCheck := true ctx := context.Background() @@ -1569,6 +1608,7 @@ func handleExecutionResult(workflowExecution WorkflowExecution) { if shutdownCheck { log.Println("BREAKING BECAUSE RESULTS IS SAME LENGTH AS ACTIONS. SHOULD CHECK ALL RESULTS FOR WHETHER THEY'RE DONE") + validateFinished(workflowExecution) shutdown(workflowExecution.ExecutionId, workflowExecution.Workflow.ID) } } @@ -1714,12 +1754,15 @@ func handleExecution(client *http.Client, req *http.Request, workflowExecution W for { handleExecutionResult(workflowExecution) - fullUrl := fmt.Sprintf("%s/api/v1/workflows/%s/executions/%s/abort", baseUrl, workflowExecution.Workflow.ID, workflowExecution.ExecutionId) + //fullUrl := fmt.Sprintf("%s/api/v1/workflows/%s/executions/%s/abort", baseUrl, workflowExecution.Workflow.ID, workflowExecution.ExecutionId) + fullUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl) + log.Printf("URL: %s", fullUrl) req, err := http.NewRequest( "POST", fullUrl, bytes.NewBuffer([]byte(data)), ) + newresp, err := topClient.Do(req) if err != nil { log.Printf("[ERROR] Failed making request: %s", err) @@ -1735,7 +1778,8 @@ func handleExecution(client *http.Client, req *http.Request, workflowExecution W } if newresp.StatusCode != 200 { - log.Printf("[ERROR] Bad statuscode: %s\nStatusCode: %d", string(body), newresp.StatusCode) + log.Printf("[ERROR] Bad statuscode: %d, %s", newresp.StatusCode, string(body)) + //shutdown(workflowExecution.ExecutionId, workflowExecution.Workflow.ID) time.Sleep(time.Duration(sleepTime) * time.Second) continue } @@ -1898,7 +1942,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { return } - log.Printf("Got result: %s", string(body)) + //log.Printf("Got result: %s", string(body)) var actionResult ActionResult err = json.Unmarshal(body, &actionResult) if err != nil { @@ -2055,7 +2099,7 @@ func findChildNodes(workflowExecution WorkflowExecution, nodeId string) []string // 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 ActionResult, resp http.ResponseWriter) { - log.Printf("IN WORKFLOWEXECUTION SUB!") + //log.Printf("IN WORKFLOWEXECUTION SUB!") // Should start a tx for the execution here workflowExecution, err := getWorkflowExecution(ctx, workflowExecutionId) if err != nil { @@ -2144,8 +2188,14 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } if !sourceNodeFound { - log.Printf("Not setting node %s to SKIPPED", nodeId) + // FIXME: Shouldn't add skip for child nodes of these nodes. Check if this node is parent of upcoming nodes. + log.Printf("\n\n NOT setting node %s to SKIPPED", nodeId) skipNodeAdd = true + + if !arrayContains(visited, nodeId) && !arrayContains(executed, nodeId) { + nextActions = append(nextActions, nodeId) + log.Printf("SHOULD EXECUTE NODE %s. Next actions: %s", nodeId, nextActions) + } break } } @@ -2163,6 +2213,11 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } newResults = append(newResults, newResult) + } else { + //log.Printf("\n\nNOT adding %s as skipaction - should add to execute?", nodeId) + //var visited []string + //var executed []string + //var nextActions []string } } } @@ -2333,6 +2388,7 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl if workflowExecution.LastNode == "" { workflowExecution.LastNode = actionResult.Action.ID } + } } @@ -2395,50 +2451,8 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl log.Printf("Skipping setexec with status %s", workflowExecution.Status) } - //ExecutionId - // Transactions: https://cloud.google.com/datastore/docs/concepts/transactions#datastore-datastore-transactional-update-go - // Prevents timing issues - //if _, err := tx.Put(key, workflowExecution); err != nil { - // log.Printf("[ERROR] tx.Put error: %v", err) - // err = tx.Rollback() - // if err != nil { - // log.Printf("[ERROR] Rollback error (3): %s", err) - // } - - // resp.WriteHeader(401) - // resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed setting workflowexecution actionresult: %s"}`, err))) - // return - //} - - //if _, err = tx.Commit(); err != nil { - // err = tx.Rollback() - // if err != nil { - // log.Printf("[ERROR] Rollback error expected ? (1): %s", err) - // } - - // if attempts >= 7 { - // log.Printf("[ERROR] QUITTING: tx.Commit %d: %v", attempts, err) - - // workflowExecution.Status = "ABORTED" - // setWorkflowExecution(ctx, *workflowExecution, true) - - // resp.WriteHeader(401) - // resp.Write([]byte(`{"success": false}`)) - // return - // } - - // if attempts > 3 { - // log.Printf("[WARNING] tx.Commit %d: %v", attempts, err) - // } - - // attempts += 1 - // runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp) - // return - //} else { - // //if grpc.Code(err) == codes.Aborted { - // // return nil, ErrConcurrentTransaction - // //} - // //t.id = nil // mark the transaction as expired + //if newExecutions && len(nextActions) > 0 { + // handleExecutionResult(*workflowExecution) //} resp.WriteHeader(200) @@ -2446,35 +2460,122 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } func getWorkflowExecution(ctx context.Context, id string) (*WorkflowExecution, error) { - log.Printf("IN GET WORKFLOW EXEC!") + //log.Printf("IN GET WORKFLOW EXEC!") cacheKey := fmt.Sprintf("workflowexecution-%s", id) if value, found := requestCache.Get(cacheKey); found { parsedValue := value.(*WorkflowExecution) //log.Printf("Found execution for id %s with %d results", parsedValue.ExecutionId, len(parsedValue.Results)) + + //validateFinished(*parsedValue) return parsedValue, nil } return &WorkflowExecution{}, errors.New("No workflowexecution defined yet") } +func validateFinished(workflowExecution WorkflowExecution) { + log.Printf("Status: %s, Actions: %d, Extra: %d, Results: %d\n", 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("[FINISHED] Should send full result to %s", baseUrl) + + //data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization) + data, err := json.Marshal(workflowExecution) + if err != nil { + log.Printf("[ERROR] Failed to unmarshal data for backend") + shutdown(workflowExecution.ExecutionId, "") + } + + fullUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl) + req, err := http.NewRequest( + "POST", + fullUrl, + bytes.NewBuffer([]byte(data)), + ) + + if err != nil { + log.Printf("[ERROR] Failed creating finishing request: %s", err) + shutdown(workflowExecution.ExecutionId, "") + } + + newresp, err := topClient.Do(req) + if err != nil { + log.Printf("[ERROR] Error running finishing request: %s", err) + shutdown(workflowExecution.ExecutionId, "") + } + + body, err := ioutil.ReadAll(newresp.Body) + log.Printf("BACKEND STATUS: %d", newresp.StatusCode) + if err != nil { + log.Printf("[ERROR] Failed reading body: %s", err) + } else { + log.Printf("NEWRESP: %s", string(body)) + } + } +} + +func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { + body, err := ioutil.ReadAll(request.Body) + if err != nil { + log.Println("Failed reading body for stream result queue") + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + var actionResult ActionResult + err = json.Unmarshal(body, &actionResult) + if err != nil { + log.Printf("Failed ActionResult unmarshaling: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + ctx := context.Background() + workflowExecution, err := getWorkflowExecution(ctx, actionResult.ExecutionId) + if err != nil { + //log.Printf("Failed getting execution (streamresult) %s: %s", actionResult.ExecutionId, err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key or execution_id might not exist."}`))) + return + } + + // Authorization is done here + if workflowExecution.Authorization != actionResult.Authorization { + log.Printf("Bad authorization key when getting stream results %s.", actionResult.ExecutionId) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key or execution_id might not exist."}`))) + return + } + + newjson, err := json.Marshal(workflowExecution) + if err != nil { + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed unpacking workflow execution"}`))) + return + } + + resp.WriteHeader(200) + resp.Write(newjson) + +} + func setWorkflowExecution(ctx context.Context, workflowExecution WorkflowExecution, dbSave bool) error { - log.Printf("IN SET WORKFLOW EXEC!") + //log.Printf("IN SET WORKFLOW EXEC!") //log.Printf("\n\n\nRESULT: %s\n\n\n", workflowExecution.Status) if len(workflowExecution.ExecutionId) == 0 { log.Printf("Workflowexeciton executionId can't be empty.") return errors.New("ExecutionId can't be empty.") } - startAction = workflowExecution.Start - if len(startAction) == 0 { - log.Printf("Didn't find execution start action. Setting it to workflow start action.") - startAction = workflowExecution.Workflow.Start - } - - handleExecutionResult(workflowExecution) - cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId) requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration) + + handleExecutionResult(workflowExecution) + validateFinished(workflowExecution) return nil } @@ -2495,22 +2596,19 @@ func GetLocalIP() string { return "" } -func runWebserver() { - //hostname, err := os.Hostname() - //if err != nil { - // log.Printf("Hostname error: %s", hostname) - // return - //} - +func webserverSetup() { hostname := GetLocalIP() log.Printf("\nStarting webserver on port 5001 with hostname: %s\n", hostname) log.Printf("OLD HOSTNAME: %s", appCallbackUrl) appCallbackUrl = fmt.Sprintf("http://%s:5001", hostname) log.Printf("NEW HOSTNAME: %s", appCallbackUrl) +} +func runWebserver() { r := mux.NewRouter() r.HandleFunc("/api/v1/streams", handleWorkflowQueue).Methods("POST") + r.HandleFunc("/api/v1/streams/results", handleGetStreamResults).Methods("POST", "OPTIONS") http.Handle("/", r) log.Fatal(http.ListenAndServe(":5001", nil)) @@ -2582,7 +2680,6 @@ func main() { topClient = client firstRequest := true - environments := []string{} for { // Because of this, it always has updated data. // Removed request requirement from app_sdk @@ -2601,7 +2698,7 @@ func main() { } if newresp.StatusCode != 200 { - log.Printf("[ERROR] %s\nStatusCode: %d", string(body), newresp.StatusCode) + log.Printf("[ERROR] %s\nStatusCode (1): %d", string(body), newresp.StatusCode) time.Sleep(time.Duration(sleepTime) * time.Second) continue } @@ -2634,16 +2731,25 @@ func main() { } } - log.Printf("Environments: %s", environments) + log.Printf("Environments: %s. 1 = webserver, 0 or >1 = default", environments) if len(environments) == 1 { + webserverSetup() err := executionInit(workflowExecution) if err != nil { log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err) shutdown(workflowExecution.ExecutionId, workflowExecution.Workflow.ID) } - handleExecutionResult(workflowExecution) + go func() { + time.Sleep(time.Duration(1)) + handleExecutionResult(workflowExecution) + }() + runWebserver() + //log.Printf("Before wait") + //wg := sync.WaitGroup{} + //wg.Add(1) + //wg.Wait() } }