#345: Added auto-replay of the entire WORKFLOW IF it hasn't finished in time.

This commit is contained in:
frikky
2021-12-20 22:54:06 +01:00
parent d74fab7bf8
commit 9e8cc0e91e
6 changed files with 65 additions and 30 deletions
+1 -1
View File
@@ -152,7 +152,7 @@ func cleanupExistingNodes(ctx context.Context) error {
//log.Printf("\n\nFound %d contaienrs", len(services))
for _, service := range services {
log.Printf("[INFO] Service: %#v", service.Spec.Annotations.Name)
//log.Printf("[INFO] Service: %#v", service.Spec.Annotations.Name)
//portFound := false
//for _, endpoint := range service.Spec.EndpointSpec.Ports {
+1 -1
View File
@@ -1,5 +1,5 @@
NAME=shuffle-worker
VERSION=0.9.42
VERSION=0.9.43
echo "Running docker build with $NAME:$VERSION"
#CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o worker.bin .
+31 -23
View File
@@ -55,6 +55,7 @@ var requestCache *cache.Cache
var topClient *http.Client
var data string
var requestsSent = 0
var appsInitialized = false
var hostname string
@@ -116,7 +117,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
}
// Might not be necessary because of cleanupEnv hostconfig autoremoval
if cleanupEnv == "true" && len(containerIds) > 0 {
if cleanupEnv == "true" && len(containerIds) > 0 && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" {
/*
ctx := context.Background()
dockercli, err := dockerclient.NewEnvClient()
@@ -152,7 +153,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
path += fmt.Sprintf("&env=%s", url.QueryEscape(environment))
}
//fmt.Println(url.QueryEscape(query))
//fmt.Printf(url.QueryEscape(query))
abortUrl += path
log.Printf("[DEBUG][%s] Abort URL: %s", workflowExecution.ExecutionId, abortUrl)
@@ -163,7 +164,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
)
if err != nil {
log.Println("[INFO][%s] Failed building request: %s", workflowExecution.ExecutionId, err)
log.Printf("[INFO][%s] Failed building request: %s", workflowExecution.ExecutionId, err)
}
// FIXME: Add an API call to the backend
@@ -255,7 +256,6 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
// FIXME: Not caring if it's ok or not. Just continuing
// This is working as intended, just designed to download an updated
// image on every Orborus/new worker restart.
downloadedImages = append(downloadedImages, image)
// Running as coroutine for eventual completeness
//go downloadDockerImageBackend(&http.Client{}, image)
@@ -713,7 +713,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if isSkipped {
//log.Printf("Skipping %s as all parents are done", item.Action.Label)
if !arrayContains(visited, item.Action.ID) {
log.Printf("[INFO][%s] Adding visited (1): %s\n", workflowExecution.ExecutionId, item.Action.Label)
log.Printf("[INFO][%s] Adding visited (1): %s", workflowExecution.ExecutionId, item.Action.Label)
visited = append(visited, item.Action.ID)
}
} else {
@@ -722,7 +722,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
} else {
if item.Status == "FINISHED" {
log.Printf("[INFO][%s] Adding visited (2): %s\n", workflowExecution.ExecutionId, item.Action.Label)
log.Printf("[INFO][%s] Adding visited (2): %s", workflowExecution.ExecutionId, item.Action.Label)
visited = append(visited, item.Action.ID)
}
}
@@ -854,7 +854,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if err != nil {
log.Printf("[DEBUG][%s] Error in skipme for %s: %s", workflowExecution.ExecutionId, action.Label, err)
} else {
log.Printf("[INFO][%s] Adding visited (4): %s\n", workflowExecution.ExecutionId, action.Label)
log.Printf("[INFO][%s] Adding visited (4): %s", workflowExecution.ExecutionId, action.Label)
visited = append(visited, action.ID)
executed = append(executed, action.ID)
@@ -1337,7 +1337,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
}
log.Printf("[INFO][%s] Adding visited (3): %s\n. Actions: %d, Results: %d", workflowExecution.ExecutionId, action.Label, len(workflowExecution.Workflow.Actions), len(workflowExecution.Results))
log.Printf("[INFO][%s] Adding visited (3): %s (%s). Actions: %d, Results: %d", workflowExecution.ExecutionId, action.Label, action.ID, len(workflowExecution.Workflow.Actions), len(workflowExecution.Results))
visited = append(visited, action.ID)
executed = append(executed, action.ID)
@@ -1347,8 +1347,8 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
//log.Printf("EXECUTED: %#v", executed)
}
//log.Println(nextAction)
//log.Println(startAction, children[startAction])
//log.Printf(nextAction)
//log.Printf(startAction, children[startAction])
// FIXME - new request here
// FIXME - clean up stopped (remove) containers with this execution id
@@ -1371,7 +1371,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
if shutdownCheck {
log.Println("[INFO][%s] BREAKING BECAUSE RESULTS IS SAME LENGTH AS ACTIONS. SHOULD CHECK ALL RESULTS FOR WHETHER THEY'RE DONE", workflowExecution.ExecutionId)
log.Printf("[INFO][%s] BREAKING BECAUSE RESULTS IS SAME LENGTH AS ACTIONS. SHOULD CHECK ALL RESULTS FOR WHETHER THEY'RE DONE", workflowExecution.ExecutionId)
validateFinished(workflowExecution)
log.Printf("[DEBUG][%s] Shutting down (17)", workflowExecution.ExecutionId)
shutdown(workflowExecution, "", "", true)
@@ -1904,7 +1904,7 @@ 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.Println("[WARNING] (3) Failed reading body for workflowqueue")
log.Printf("[WARNING] (3) Failed reading body for workflowqueue")
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
return
@@ -2144,7 +2144,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
//log.Printf("[DEBUG] Got stream result")
body, err := ioutil.ReadAll(request.Body)
if err != nil {
log.Println("Failed reading body for stream result queue")
log.Printf("Failed reading body for stream result queue")
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
return
@@ -2346,6 +2346,9 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
func downloadDockerImageBackend(client *http.Client, imageName string) error {
log.Printf("[DEBUG] Trying to download image %s from backend as it doesn't exist", imageName)
downloadedImages = append(downloadedImages, imageName)
data := fmt.Sprintf(`{"name": "%s"}`, imageName)
dockerImgUrl := fmt.Sprintf("%s/api/v1/get_docker_image", baseUrl)
@@ -2652,8 +2655,11 @@ func findAppInfo(image, name string) (int, error) {
}
exposedPort = highest
log.Printf("[DEBUG] Waiting 10 seconds before moving on to let app start")
time.Sleep(time.Duration(10) * time.Second)
if appsInitialized {
log.Printf("[DEBUG] Waiting 10 seconds before moving on to let app start")
time.Sleep(time.Duration(10) * time.Second)
}
}
return exposedPort, nil
@@ -2785,7 +2791,7 @@ func sendAppRequest(incomingUrl, appName string, port int, action shuffle.Action
// Function to auto-deploy certain apps if "run" is set
// Has some issues with loading when running multiple workers and such.
func baseDeploy() {
return
//return
cli, err := dockerclient.NewEnvClient()
if err != nil {
@@ -2832,12 +2838,14 @@ func baseDeploy() {
//deployApp(cli, value, identifier, env, workflowExecution, action)
log.Printf("[DEBUG] Deploying app with identifier %s to ensure basic apps are available from the get-go", identifier)
go deployApp(cli, value, identifier, env, workflowExecution, action)
deployApp(cli, value, identifier, env, workflowExecution, action)
//err := deployApp(cli, value, identifier, env, workflowExecution, action)
//if err != nil {
// log.Printf("[DEBUG] Failed deploying app %s: %s", value, err)
//}
}
appsInitialized = true
}
// Initial loop etc
@@ -2898,8 +2906,8 @@ func main() {
//var autoDeploy = []string{"frikky/shuffle:shuffle-subflow_1.0.0", "frikky/shuffle:http_1.1.0", "frikky/shuffle:shuffle-tools_1.1.0", "frikky/shuffle:testing_1.0.0"}
//go baseDeploy()
baseDeploy()
go baseDeploy()
//baseDeploy()
listener := webserverSetup(workflowExecution)
runWebserver(listener)
@@ -2931,13 +2939,13 @@ func main() {
ExecutionId: executionId,
}
if len(authorization) == 0 {
log.Println("[INFO] No AUTHORIZATION key set in env")
log.Printf("[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("[INFO] No EXECUTIONID key set in env")
log.Printf("[DEBUG] Shutting down (28)")
shutdown(workflowExecution, "", "", false)
}
@@ -2951,7 +2959,7 @@ func main() {
)
if err != nil {
log.Println("[ERROR] Failed making request builder for backend")
log.Printf("[ERROR] Failed making request builder for backend")
log.Printf("[DEBUG] Shutting down (29)")
shutdown(workflowExecution, "", "", true)
}
@@ -3099,7 +3107,7 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
body, err := ioutil.ReadAll(request.Body)
if err != nil {
log.Println("[WARNING] Failed reading body for stream result queue")
log.Printf("[WARNING] Failed reading body for stream result queue")
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
return