#396: Fixed a bug where workflows are randomly aborted

This commit is contained in:
frikky
2021-06-07 06:47:24 +02:00
parent d6f1b0e3ae
commit 21279e7ea6
+50 -11
View File
@@ -93,7 +93,7 @@ func init() {
// removes every container except itself (worker)
func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason string, handleResultSend bool) {
log.Printf("[INFO] Shutdown (%s) started with reason %s", workflowExecution.Status, reason)
log.Printf("[INFO] Shutdown (%s) started with reason %s. Result amount: %d. ResultsSent: %d, Send result: %#v", workflowExecution.Status, reason, len(workflowExecution.Results), requestsSent, handleResultSend)
//reason := "Error in execution"
sleepDuration := 1
@@ -101,9 +101,9 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
shutdownData, err := json.Marshal(workflowExecution)
if err == nil {
sendResult(workflowExecution, shutdownData)
log.Printf("[WARNING] Sent shutdown update")
log.Printf("[WARNING] Sent shutdown update with %d results and result value %s", len(workflowExecution.Results), reason)
} else {
log.Printf("[WARNING] DIDNT send update")
log.Printf("[WARNING] Failed to send update: %s", err)
}
time.Sleep(time.Duration(sleepDuration) * time.Second)
@@ -554,7 +554,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
for _, subresult := range workflowExecution.Results {
if subresult.Action.ID == branch.SourceID {
if subresult.Status != "SKIPPED" && subresult.Status != "FAILURE" {
log.Printf("\n\n\nSUBRESULT PARENT STATUS: %s\n\n\n", subresult.Status)
//log.Printf("\n\n\nSUBRESULT PARENT STATUS: %s\n\n\n", subresult.Status)
isSkipped = false
break
@@ -617,7 +617,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
if exit && len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions) {
log.Printf("Shutting down.")
log.Printf("[DEBUG] Shutting down (1)")
shutdown(workflowExecution, "", "", true)
}
@@ -976,6 +976,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
err = deployApp(dockercli, images[0], identifier, env, workflowExecution)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (2)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -983,12 +984,14 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
err = deployApp(dockercli, image, identifier, env, workflowExecution)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (3)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
log.Printf("[WARNING] Failed CLEANUP execution. Downloading image remotely.")
reader, err := dockercli.ImagePull(context.Background(), image, pullOptions)
if err != nil {
log.Printf("[ERROR] Failed getting %s. Couldn't be find locally, AND is missing.", image)
log.Printf("[DEBUG] Shutting down (4)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -996,10 +999,12 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
_, err = io.Copy(buildBuf, reader)
if err != nil && !strings.Contains(fmt.Sprintf("%s", err.Error()), "Conflict. The container name") {
log.Printf("[ERROR] Error in IO copy: %s", err)
log.Printf("[DEBUG] Shutting down (5)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
} else {
if strings.Contains(buildBuf.String(), "errorDetail") {
log.Printf("[ERROR] Docker build:\n%s\nERROR ABOVE: Trying to pull tags from: %s", buildBuf.String(), image)
log.Printf("[DEBUG] Shutting down (6)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1011,12 +1016,14 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
log.Printf("[ERROR] Failed deploying image for the FOURTH time. Aborting if the image doesn't exist")
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (7)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
if strings.Contains(err.Error(), "No such image") {
//log.Printf("[WARNING] Failed deploying %s from image %s: %s", identifier, image, err)
log.Printf("[ERROR] Image doesn't exist. Shutting down")
log.Printf("[DEBUG] Shutting down (8)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
}
@@ -1027,6 +1034,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
err = deployApp(dockercli, images[0], identifier, env, workflowExecution)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (9)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1036,6 +1044,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
err = deployApp(dockercli, image, identifier, env, workflowExecution)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (10)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1043,6 +1052,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
err = deployApp(dockercli, image, identifier, env, workflowExecution)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (11)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1050,6 +1060,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
reader, err := dockercli.ImagePull(context.Background(), image, pullOptions)
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
log.Printf("[ERROR] Failed getting %s. The couldn't be find locally, AND is missing.", image)
log.Printf("[DEBUG] Shutting down (12)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1057,10 +1068,12 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
_, err = io.Copy(buildBuf, reader)
if err != nil {
log.Printf("[ERROR] Error in IO copy: %s", err)
log.Printf("[DEBUG] Shutting down (13)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
} else {
if strings.Contains(buildBuf.String(), "errorDetail") {
log.Printf("[ERROR] Docker build:\n%s\nERROR ABOVE: Trying to pull tags from: %s", buildBuf.String(), image)
log.Printf("[DEBUG] Shutting down (14)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
@@ -1071,12 +1084,14 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") {
log.Printf("[ERROR] Failed deploying image for the FOURTH time. Aborting if the image doesn't exist")
if strings.Contains(err.Error(), "exited prematurely") {
log.Printf("[DEBUG] Shutting down (15)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
if strings.Contains(err.Error(), "No such image") {
//log.Printf("[WARNING] Failed deploying %s from image %s: %s", identifier, image, err)
log.Printf("[ERROR] Image doesn't exist. Shutting down")
log.Printf("[DEBUG] Shutting down (16)")
shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true)
}
}
@@ -1117,6 +1132,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if shutdownCheck {
log.Println("[INFO] BREAKING BECAUSE RESULTS IS SAME LENGTH AS ACTIONS. SHOULD CHECK ALL RESULTS FOR WHETHER THEY'RE DONE")
validateFinished(workflowExecution)
log.Printf("[DEBUG] Shutting down (17)")
shutdown(workflowExecution, "", "", true)
}
}
@@ -1255,6 +1271,7 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
err := executionInit(workflowExecution)
if err != nil {
log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err)
log.Printf("[DEBUG] Shutting down (18)")
shutdown(workflowExecution, "", "", true)
}
@@ -1291,7 +1308,8 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
log.Printf("[ERROR] Bad statuscode: %d, %s", newresp.StatusCode, string(body))
if strings.Contains(string(body), "Workflowexecution is already finished") {
shutdown(workflowExecution, "", "", false)
log.Printf("[DEBUG] Shutting down (19)")
shutdown(workflowExecution, "", "", true)
}
time.Sleep(time.Duration(sleepTime) * time.Second)
@@ -1307,12 +1325,14 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
log.Printf("[INFO] Workflow %s is finished. Exiting worker.", workflowExecution.ExecutionId)
log.Printf("[DEBUG] Shutting down (20)")
shutdown(workflowExecution, "", "", true)
}
log.Printf("[INFO] Status: %s, Results: %d, actions: %d", workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
if workflowExecution.Status != "EXECUTING" {
log.Printf("[WARNING] Exiting as worker execution has status %s!", workflowExecution.Status)
log.Printf("[DEBUG] Shutting down (21)")
shutdown(workflowExecution, "", "", true)
}
@@ -1544,7 +1564,7 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
//dbSave := false
if len(results) != len(workflowExecution.Results) {
log.Printf("\n\n[WARNING] 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.\n\n", 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))
}
// Validating that action results hasn't changed
@@ -1619,12 +1639,14 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
if err != nil {
log.Printf("[ERROR] Failed creating finishing request: %s", err)
log.Printf("[DEBUG] Shutting down (22)")
shutdown(workflowExecution, "", "", false)
}
newresp, err := topClient.Do(req)
if err != nil {
log.Printf("[ERROR] Error running finishing request: %s", err)
log.Printf("[DEBUG] Shutting down (23)")
shutdown(workflowExecution, "", "", false)
}
@@ -1649,6 +1671,7 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) {
shutdownData, err := json.Marshal(workflowExecution)
if err != nil {
log.Printf("[ERROR] Failed to unmarshal data for backend")
log.Printf("[DEBUG] Shutting down (24)")
shutdown(workflowExecution, "", "", true)
}
@@ -1704,10 +1727,8 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
}
func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.WorkflowExecution, dbSave bool) error {
//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.")
log.Printf("[INFO] Workflowexecution executionId can't be empty.")
return errors.New("ExecutionId can't be empty.")
}
@@ -1717,8 +1738,18 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
handleExecutionResult(workflowExecution)
validateFinished(workflowExecution)
// FIXME: Should this shutdown OR send the result?
// The worker may not be running the backend hmm
if dbSave {
shutdown(workflowExecution, "", "", false)
if workflowExecution.ExecutionSource == "default" {
log.Printf("[DEBUG] Shutting down (25)")
shutdown(workflowExecution, "", "", true)
//log.Printf("[INFO] Not sending backend info since source is default")
//return
} else {
log.Printf("[DEBUG] NOT shutting down with dbSave (%s)", workflowExecution.ExecutionSource)
}
}
return nil
@@ -1761,6 +1792,7 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
listener, err := getAvailablePort()
if err != nil {
log.Printf("Failed to created listener: %s", err)
log.Printf("[DEBUG] Shutting down (26)")
shutdown(workflowExecution, "", "", true)
}
port := listener.Addr().(*net.TCPAddr).Port
@@ -1918,11 +1950,13 @@ func main() {
}
if len(authorization) == 0 {
log.Println("[INFO] No AUTHORIZATION key set in env")
log.Printf("[DEBUG] Shutting down (27)")
shutdown(workflowExecution, "", "", false)
}
if len(executionId) == 0 {
log.Println("[INFO] No EXECUTIONID key set in env")
log.Printf("[DEBUG] Shutting down (28)")
shutdown(workflowExecution, "", "", false)
}
@@ -1936,6 +1970,7 @@ func main() {
if err != nil {
log.Println("[ERROR] Failed making request builder for backend")
log.Printf("[DEBUG] Shutting down (29)")
shutdown(workflowExecution, "", "", true)
}
topClient = client
@@ -1999,6 +2034,7 @@ func main() {
err := executionInit(workflowExecution)
if err != nil {
log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err)
log.Printf("[DEBUG] Shutting down (30)")
shutdown(workflowExecution, "", "", true)
}
@@ -2024,6 +2060,7 @@ func main() {
if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
log.Printf("[INFO] Workflow %s is finished. Exiting worker.", workflowExecution.ExecutionId)
log.Printf("[DEBUG] Shutting down (31)")
shutdown(workflowExecution, "", "", true)
}
@@ -2032,10 +2069,12 @@ func main() {
err = handleDefaultExecution(client, req, workflowExecution)
if err != nil {
log.Printf("[INFO] Workflow %s is finished: %s", workflowExecution.ExecutionId, err)
log.Printf("[DEBUG] Shutting down (32)")
shutdown(workflowExecution, "", "", true)
}
} else {
log.Printf("[INFO] Workflow %s has status %s. Exiting worker.", workflowExecution.ExecutionId, workflowExecution.Status)
log.Printf("[DEBUG] Shutting down (33)")
shutdown(workflowExecution, workflowExecution.Workflow.ID, "", true)
}