|
|
|
@@ -33,6 +33,10 @@ import (
|
|
|
|
|
"github.com/gorilla/mux"
|
|
|
|
|
"github.com/patrickmn/go-cache"
|
|
|
|
|
"github.com/satori/go.uuid"
|
|
|
|
|
|
|
|
|
|
// No necessary outside shared
|
|
|
|
|
"cloud.google.com/go/datastore"
|
|
|
|
|
"cloud.google.com/go/storage"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
// This is getting out of hand :)
|
|
|
|
@@ -50,17 +54,19 @@ 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
|
|
|
|
|
var executed []string
|
|
|
|
|
var nextActions []string
|
|
|
|
|
var containerIds []string
|
|
|
|
|
var extra int
|
|
|
|
|
var startAction string
|
|
|
|
|
*/
|
|
|
|
|
var results []shuffle.ActionResult
|
|
|
|
|
var allLogs map[string]string
|
|
|
|
|
var containerIds []string
|
|
|
|
|
|
|
|
|
|
var executionRunning bool
|
|
|
|
|
|
|
|
|
@@ -137,11 +143,15 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// FIXME: Add an API call to the backend
|
|
|
|
|
authorization := os.Getenv("AUTHORIZATION")
|
|
|
|
|
if len(authorization) > 0 {
|
|
|
|
|
req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization))
|
|
|
|
|
if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" {
|
|
|
|
|
authorization := os.Getenv("AUTHORIZATION")
|
|
|
|
|
if len(authorization) > 0 {
|
|
|
|
|
req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization))
|
|
|
|
|
} else {
|
|
|
|
|
log.Printf("[ERROR] No authorization specified for abort")
|
|
|
|
|
}
|
|
|
|
|
} else {
|
|
|
|
|
log.Printf("[ERROR] No authorization specified for abort")
|
|
|
|
|
req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", workflowExecution.Authorization))
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
req.Header.Add("Content-Type", "application/json")
|
|
|
|
@@ -182,19 +192,22 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
|
|
|
|
|
os.Exit(3)
|
|
|
|
|
} else {
|
|
|
|
|
log.Printf("\n\n[DEBUG] Sending result and resetting values (K8s & Swarm).\n\n")
|
|
|
|
|
environments = []string{}
|
|
|
|
|
parents = map[string][]string{}
|
|
|
|
|
children = map[string][]string{}
|
|
|
|
|
visited = []string{}
|
|
|
|
|
executed = []string{}
|
|
|
|
|
nextActions = []string{}
|
|
|
|
|
containerIds = []string{}
|
|
|
|
|
extra = 0
|
|
|
|
|
startAction = ""
|
|
|
|
|
results = []shuffle.ActionResult{}
|
|
|
|
|
allLogs = map[string]string{}
|
|
|
|
|
requestsSent = 0
|
|
|
|
|
//UpdateExecutionVariables(ctx, workflowExecution.ExecutionId, startAction, children, parents, visited, executed, nextActions, environments, extra)
|
|
|
|
|
|
|
|
|
|
/*
|
|
|
|
|
environments = []string{}
|
|
|
|
|
parents = map[string][]string{}
|
|
|
|
|
children = map[string][]string{}
|
|
|
|
|
visited = []string{}
|
|
|
|
|
executed = []string{}
|
|
|
|
|
nextActions = []string{}
|
|
|
|
|
containerIds = []string{}
|
|
|
|
|
extra = 0
|
|
|
|
|
startAction = ""
|
|
|
|
|
results = []shuffle.ActionResult{}
|
|
|
|
|
allLogs = map[string]string{}
|
|
|
|
|
*/
|
|
|
|
|
requestsSent = 0
|
|
|
|
|
executionRunning = false
|
|
|
|
|
}
|
|
|
|
|
//cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId)
|
|
|
|
@@ -563,7 +576,11 @@ func removeIndex(s []string, i int) []string {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
log.Printf("[INFO] Inside execution results with %d / %d results", len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
|
|
|
|
|
|
|
|
|
|
if len(startAction) == 0 {
|
|
|
|
|
startAction = workflowExecution.Start
|
|
|
|
|
if len(startAction) == 0 {
|
|
|
|
@@ -1275,12 +1292,14 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func executionInit(workflowExecution shuffle.WorkflowExecution) error {
|
|
|
|
|
parents = map[string][]string{}
|
|
|
|
|
children = map[string][]string{}
|
|
|
|
|
parents := map[string][]string{}
|
|
|
|
|
children := map[string][]string{}
|
|
|
|
|
nextActions := []string{}
|
|
|
|
|
extra := 0
|
|
|
|
|
|
|
|
|
|
results = workflowExecution.Results
|
|
|
|
|
|
|
|
|
|
startAction = workflowExecution.Start
|
|
|
|
|
startAction := workflowExecution.Start
|
|
|
|
|
log.Printf("[INFO] STARTACTION: %s", startAction)
|
|
|
|
|
if len(startAction) == 0 {
|
|
|
|
|
log.Printf("[INFO] Didn't find execution start action. Setting it to workflow start action.")
|
|
|
|
@@ -1375,7 +1394,7 @@ func executionInit(workflowExecution shuffle.WorkflowExecution) error {
|
|
|
|
|
pullOptions := types.ImagePullOptions{}
|
|
|
|
|
_ = pullOptions
|
|
|
|
|
for _, image := range onpremApps {
|
|
|
|
|
log.Printf("[INFO] Image: %s", image)
|
|
|
|
|
//log.Printf("[INFO] Image: %s", image)
|
|
|
|
|
// Kind of gambling that the image exists.
|
|
|
|
|
if strings.Contains(image, " ") {
|
|
|
|
|
image = strings.ReplaceAll(image, " ", "-")
|
|
|
|
@@ -1394,12 +1413,41 @@ func executionInit(workflowExecution shuffle.WorkflowExecution) error {
|
|
|
|
|
//log.Printf("Successfully downloaded and built %s", image)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
|
|
|
|
|
visited := []string{}
|
|
|
|
|
executed := []string{}
|
|
|
|
|
environments := []string{}
|
|
|
|
|
for _, action := range workflowExecution.Workflow.Actions {
|
|
|
|
|
found := false
|
|
|
|
|
|
|
|
|
|
for _, environment := range environments {
|
|
|
|
|
if action.Environment == environment {
|
|
|
|
|
found = true
|
|
|
|
|
break
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if !found {
|
|
|
|
|
environments = append(environments, action.Environment)
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
//var visited []string
|
|
|
|
|
//var executed []string
|
|
|
|
|
err := shuffle.UpdateExecutionVariables(ctx, workflowExecution.ExecutionId, startAction, children, parents, visited, executed, nextActions, environments, extra)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("\n\n[ERROR] Failed to update exec variables for execution %s: %s\n\n", workflowExecution.ExecutionId, err)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func handleDefaultExecution(client *http.Client, req *http.Request, workflowExecution shuffle.WorkflowExecution) error {
|
|
|
|
|
// if no onprem runs (shouldn't happen, but extra check), exit
|
|
|
|
|
// if there are some, load the images ASAP for the app
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
//startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
|
|
|
|
startAction, extra, _, _, _, _, _, _ := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
err := executionInit(workflowExecution)
|
|
|
|
|
if err != nil {
|
|
|
|
@@ -1410,7 +1458,6 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
|
|
|
|
|
|
|
|
|
|
log.Printf("[DEBUG] DEFAULT EXECUTION Startaction: %s", startAction)
|
|
|
|
|
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
setWorkflowExecution(ctx, workflowExecution, false)
|
|
|
|
|
|
|
|
|
|
streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl)
|
|
|
|
@@ -1649,7 +1696,7 @@ func runTestExecution(client *http.Client, workflowId, apikey string) (string, s
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
log.Printf("[DEBUG] Got stream workflow queue")
|
|
|
|
|
//log.Printf("[DEBUG] Got stream workflow queue")
|
|
|
|
|
body, err := ioutil.ReadAll(request.Body)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("(3) Failed reading body for workflowqueue")
|
|
|
|
@@ -1693,7 +1740,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
if workflowExecution.Status == "FINISHED" {
|
|
|
|
|
log.Printf("[DEBUG] Workflowexecution is already FINISHED. No further action can be taken")
|
|
|
|
|
resp.WriteHeader(401)
|
|
|
|
|
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Workflowexecution is already finished because it has status %s"}`, workflowExecution.LastNode, workflowExecution.Status)))
|
|
|
|
|
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Workflowexecution is already finished because it has status %s. Lastnode: %s"}`, workflowExecution.Status, workflowExecution.LastNode)))
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -1863,6 +1910,10 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func validateFinished(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
ctx := context.Background()
|
|
|
|
|
//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))
|
|
|
|
|
|
|
|
|
|
//if len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions)+extra {
|
|
|
|
@@ -1883,7 +1934,7 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
log.Printf("[DEBUG] Got stream result")
|
|
|
|
|
//log.Printf("[DEBUG] Got stream result")
|
|
|
|
|
body, err := ioutil.ReadAll(request.Body)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("Failed reading body for stream result queue")
|
|
|
|
@@ -1998,7 +2049,7 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
|
|
|
|
|
return listener
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("OLD HOSTNAME: %s", appCallbackUrl)
|
|
|
|
|
log.Printf("[DEBUG] OLD HOSTNAME: %s", appCallbackUrl)
|
|
|
|
|
if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" {
|
|
|
|
|
log.Printf("\n\nStarting webserver on port %d with hostname: %s\n\n", baseport, hostname)
|
|
|
|
|
appCallbackUrl = fmt.Sprintf("http://%s:%d", hostname, baseport)
|
|
|
|
@@ -2091,7 +2142,7 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error {
|
|
|
|
|
|
|
|
|
|
imageLoadResponse, err := dockercli.ImageLoad(context.Background(), tar, true)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Error loading: %s", err)
|
|
|
|
|
log.Printf("[ERROR] Error loading images: %s", err)
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -2232,9 +2283,28 @@ func findAppInfo(image, name string) (int, error) {
|
|
|
|
|
serviceListOptions,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
// Basic self-correction
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Unable to list services: %s", err)
|
|
|
|
|
return -1, err
|
|
|
|
|
log.Printf("[ERROR] Unable to list services: %s (may continue anyway?)", err)
|
|
|
|
|
if strings.Contains(fmt.Sprintf("%s", err), "is too new") {
|
|
|
|
|
// Static for some reason
|
|
|
|
|
defaultVersion := "1.40"
|
|
|
|
|
dockerApiVersion = defaultVersion
|
|
|
|
|
os.Setenv("DOCKER_API_VERSION", defaultVersion)
|
|
|
|
|
log.Printf("[DEBUG] Setting Docker API to %s default and retrying listing requests", defaultVersion)
|
|
|
|
|
} else {
|
|
|
|
|
return -1, err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
services, err = dockercli.ServiceList(
|
|
|
|
|
context.Background(),
|
|
|
|
|
serviceListOptions,
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Unable to list services (2): %s", err)
|
|
|
|
|
return -1, err
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
for _, service := range services {
|
|
|
|
@@ -2393,6 +2463,13 @@ func main() {
|
|
|
|
|
}
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
_, err := shuffle.RunInit(datastore.Client{}, storage.Client{}, "", "", false, "")
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Failed to run worker init: %s", err)
|
|
|
|
|
} else {
|
|
|
|
|
log.Printf("[DEBUG] Ran init for worker to set up cache system. Docker version: %s", dockerApiVersion)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("[INFO] Setting up worker environment")
|
|
|
|
|
sleepTime := 5
|
|
|
|
|
client := &http.Client{
|
|
|
|
@@ -2418,7 +2495,7 @@ func main() {
|
|
|
|
|
timezone = "Europe/Amsterdam"
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("[INFO] Running with timezone %s", timezone)
|
|
|
|
|
log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, os.Getenv("SHUFFLE_SWARM_CONFIG"))
|
|
|
|
|
if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" {
|
|
|
|
|
workflowExecution := shuffle.WorkflowExecution{}
|
|
|
|
|
listener := webserverSetup(workflowExecution)
|
|
|
|
@@ -2489,6 +2566,7 @@ func main() {
|
|
|
|
|
topClient = client
|
|
|
|
|
|
|
|
|
|
firstRequest := true
|
|
|
|
|
environments := []string{}
|
|
|
|
|
for {
|
|
|
|
|
// Because of this, it always has updated data.
|
|
|
|
|
// Removed request requirement from app_sdk
|
|
|
|
@@ -2676,6 +2754,7 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
//if strings.ToLower(os.Getenv("SHUFFLE_PASS_APP_PROXY")) == "true" {
|
|
|
|
|
// Is it ok if these are standard? Should they be update-able after launch? Hmm
|
|
|
|
|
if len(execRequest.HTTPProxy) > 0 {
|
|
|
|
|
log.Printf("[DEBUG] Sending proxy info to child process")
|
|
|
|
|
os.Setenv("SHUFFLE_PASS_APP_PROXY", execRequest.ShufflePassProxyToApp)
|
|
|
|
@@ -2761,11 +2840,23 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
executionRunning = false
|
|
|
|
|
log.Printf("[INFO] Workflow %s is finished. Exiting worker.", workflowExecution.ExecutionId)
|
|
|
|
|
log.Printf("[DEBUG] Shutting down (20)")
|
|
|
|
|
resp.WriteHeader(401)
|
|
|
|
|
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad status %s"}`, workflowExecution.Status)))
|
|
|
|
|
|
|
|
|
|
resp.WriteHeader(200)
|
|
|
|
|
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad status for execution - already %s. Returning with 200 OK"}`, workflowExecution.Status)))
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
//ctx := context.Background()
|
|
|
|
|
//startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
extra := 0
|
|
|
|
|
for _, trigger := range workflowExecution.Workflow.Triggers {
|
|
|
|
|
//log.Printf("Appname trigger (0): %s", trigger.AppName)
|
|
|
|
|
if trigger.AppName == "User Input" || trigger.AppName == "Shuffle Workflow" {
|
|
|
|
|
extra += 1
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("[INFO] Status: %s, Results: %d, actions: %d", workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
|
|
|
|
|
if workflowExecution.Status != "EXECUTING" {
|
|
|
|
|
executionRunning = false
|
|
|
|
@@ -2787,7 +2878,10 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err)
|
|
|
|
|
log.Printf("[DEBUG] Shutting down (30)")
|
|
|
|
|
shutdown(workflowExecution, "", "", true)
|
|
|
|
|
resp.WriteHeader(401)
|
|
|
|
|
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Error in execution init: %s"}`, err)))
|
|
|
|
|
return
|
|
|
|
|
//shutdown(workflowExecution, "", "", true)
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
handleExecutionResult(workflowExecution)
|
|
|
|
|