Worker a lot optimized

This commit is contained in:
frikky
2021-01-12 17:59:14 +01:00
parent 43e228929e
commit f0a299e69a
+185 -79
View File
@@ -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()
}
}