From 607c9adddb212821decdcf720e5387f2aa22954c Mon Sep 17 00:00:00 2001 From: frikky Date: Tue, 26 Oct 2021 00:03:59 +0200 Subject: [PATCH] Started migrating Orborus and Workers to microservice and scaling architecture --- backend/go-app/walkoff.go | 24 +- functions/onprem/orborus/orborus.go | 170 ++++++++++++++- functions/onprem/worker/build.sh | 2 +- functions/onprem/worker/go.sum | 2 + functions/onprem/worker/worker.go | 327 +++++++++++++++++++++++++--- 5 files changed, 473 insertions(+), 52 deletions(-) diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 365017a0..2475e57a 100644 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -676,20 +676,10 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { } ctx := context.Background() - executionRequests, err := shuffle.GetWorkflowQueue(ctx, id) - if err != nil { - // Skipping as this comes up over and over - //log.Printf("(2) Failed reading body for workflowqueue: %s", err) - resp.WriteHeader(401) - resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) - return - } - env, err := shuffle.GetEnvironment(ctx, id, "") timeNow := time.Now().Unix() if err == nil && len(env.Id) > 0 && len(env.Name) > 0 { if time.Now().Unix() > env.Edited+60 { - log.Printf("[DEBUG] Updating env with IP %s!", request.RemoteAddr) env.RunningIp = request.RemoteAddr env.Checkin = timeNow err = shuffle.SetEnvironment(ctx, env) @@ -699,11 +689,20 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { } } + executionRequests, err := shuffle.GetWorkflowQueue(ctx, id) + if err != nil { + // Skipping as this comes up over and over + //log.Printf("(2) Failed reading body for workflowqueue: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + // Checking and updating the environment related to the first execution if len(executionRequests.Data) == 0 { executionRequests.Data = []shuffle.ExecutionRequest{} } else { - log.Printf("In workflowqueue with %d", len(executionRequests.Data)) + //log.Printf("In workflowqueue with %d", len(executionRequests.Data)) // Try again :) if len(env.Id) == 0 && len(env.Name) == 0 { @@ -726,9 +725,7 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { //resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "No env found matching %s"}`, id))) //return } else { - log.Printf("Found Env: %#v", env) if timeNow > env.Edited+60 { - log.Printf("Updating env with IP %s!", request.RemoteAddr) env.RunningIp = request.RemoteAddr env.Checkin = timeNow err = shuffle.SetEnvironment(ctx, env) @@ -743,7 +740,6 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { if len(executionRequests.Data) > 10 { executionRequests.Data = executionRequests.Data[0:9] } - log.Printf("In workflowqueue with %d (2)", len(executionRequests.Data)) } newjson, err := json.Marshal(executionRequests) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index b50e70f7..24ffb94d 100644 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -1,9 +1,11 @@ package main /* - Orborus exists to listen for new workflow executions and deploy workers. + Orborus exists to listen for new workflow executions whcih are deployed as workers. */ +// frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/frikky/shuffle-worker:nightly + import ( "github.com/shuffle/shuffle-shared" @@ -23,6 +25,8 @@ import ( "github.com/docker/docker/api/types" "github.com/docker/docker/api/types/container" + "github.com/docker/docker/api/types/mount" + "github.com/docker/docker/api/types/swarm" //"github.com/docker/docker/api/types/filters" dockerclient "github.com/docker/docker/client" "github.com/satori/go.uuid" @@ -57,6 +61,7 @@ var runningMode = strings.ToLower(os.Getenv("RUNNING_MODE")) var cleanupEnv = strings.ToLower(os.Getenv("CLEANUP")) var timezone = os.Getenv("TZ") var containerName = os.Getenv("ORBORUS_CONTAINER_NAME") +var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG") var executionIds = []string{} var dockercli *dockerclient.Client @@ -129,7 +134,7 @@ func getThisContainerId() { // Deploys the internal worker whenever something happens // https://docs.docker.com/engine/api/sdk/examples/ -func deployWorker(image string, identifier string, env []string) { +func deployWorker(image string, identifier string, env []string, executionRequest shuffle.ExecutionRequest) { // Binds is the actual "-v" volume. // Max 20% CPU every second @@ -157,6 +162,83 @@ func deployWorker(image string, identifier string, env []string) { Env: env, } + //var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG") + parsedUuid := uuid.NewV4() + if swarmConfig == "run" { + // frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/frikky/shuffle-worker:nightly + + log.Printf("[DEBUG] Deploying containers with swarm") + //containerName := fmt.Sprintf("shuffle-worker-%s", parsedUuid) + containerName := fmt.Sprintf("shuffle-workers") + serviceSpec := swarm.ServiceSpec{ + Annotations: swarm.Annotations{ + Name: containerName, + Labels: map[string]string{}, + }, + EndpointSpec: &swarm.EndpointSpec{ + Ports: []swarm.PortConfig{ + swarm.PortConfig{ + Protocol: swarm.PortConfigProtocolTCP, + PublishMode: swarm.PortConfigPublishModeIngress, + Name: "worker-port", + PublishedPort: 33333, + TargetPort: 33333, + }, + }, + }, + TaskTemplate: swarm.TaskSpec{ + Resources: &swarm.ResourceRequirements{ + Reservations: &swarm.Resources{}, + }, + ContainerSpec: &swarm.ContainerSpec{ + Image: image, + Env: []string{ + fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", os.Getenv("SHUFFLE_SWARM_CONFIG")), + }, + Mounts: []mount.Mount{ + mount.Mount{ + Source: "/var/run/docker.sock", + Target: "/var/run/docker.sock", + Type: mount.TypeBind, + }, + }, + }, + RestartPolicy: &swarm.RestartPolicy{ + Condition: swarm.RestartPolicyConditionNone, + }, + Placement: &swarm.Placement{ + MaxReplicas: 1, + }, + }, + } + + if dockerApiVersion != "" { + serviceSpec.TaskTemplate.ContainerSpec.Env = append(serviceSpec.TaskTemplate.ContainerSpec.Env, fmt.Sprintf("DOCKER_API_VERSION=%s", dockerApiVersion)) + } + + serviceOptions := types.ServiceCreateOptions{} + service, err := dockercli.ServiceCreate( + context.Background(), + serviceSpec, + serviceOptions, + ) + + if err == nil { + log.Printf("[DEBUG] Waiting 10 seconds for workers to come awake") + time.Sleep(time.Duration(10) * time.Second) + } + + log.Printf("Servicecreate request: %#v %#v", service, err) + + err = sendWorkerRequest(executionRequest) + if err != nil { + log.Printf("[ERROR] Failed worker request: %s", err) + } else { + log.Printf("[DEBUG] Started worker from request: %s - %#v - %s", containerName, service, err) + } + return + } + //log.Printf("[INFO] Identifier: %s", identifier) cont, err := dockercli.ContainerCreate( context.Background(), @@ -169,8 +251,7 @@ func deployWorker(image string, identifier string, env []string) { if err != nil { if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") { - uuid := uuid.NewV4() - identifier = fmt.Sprintf("%s-%s", identifier, uuid) + identifier = fmt.Sprintf("%s-%s", identifier, parsedUuid) log.Printf("[INFO] 2 - Identifier: %s", identifier) cont, err = dockercli.ContainerCreate( context.Background(), @@ -212,7 +293,7 @@ func deployWorker(image string, identifier string, env []string) { // return // } - // err = deployWorker(cli, workerImage, containerName, env) + // err = deployWorke(cli, workerImage, containerName, env) // if err != nil { // log.Printf("Failed executing worker %s in state %s", execution.ExecutionId, containerStatus) // return @@ -263,14 +344,16 @@ func initializeImages() { if baseimageregistry == "" { baseimageregistry = "docker.io" baseimageregistry = "ghcr.io" - log.Printf("Setting baseimageregistry") + log.Printf("[DEBUG] Setting baseimageregistry") } if baseimagename == "" { baseimagename = "frikky/shuffle" baseimagename = "frikky" - log.Printf("Setting baseimagename") + log.Printf("[DEBUG] Setting baseimagename") } + log.Printf("[DEBUG] Setting swarm config to %#v. Default is empty.", swarmConfig) + // check whether they are the same first images := []string{ fmt.Sprintf("frikky/shuffle:app_sdk"), @@ -569,6 +652,7 @@ func main() { fmt.Sprintf("CLEANUP=%s", cleanupEnv), fmt.Sprintf("TZ=%s", timezone), fmt.Sprintf("SHUFFLE_PASS_APP_PROXY=%s", os.Getenv("SHUFFLE_PASS_APP_PROXY")), + fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", os.Getenv("SHUFFLE_SWARM_CONFIG")), } //log.Printf("Running worker with proxy? %s", os.Getenv("SHUFFLE_PASS_WORKER_PROXY")) @@ -581,7 +665,7 @@ func main() { env = append(env, fmt.Sprintf("DOCKER_API_VERSION=%s", dockerApiVersion)) } - go deployWorker(workerImage, containerName, env) + go deployWorker(workerImage, containerName, env, execution) log.Printf("[INFO] ExecutionID %s was deployed and to be removed from queue.", execution.ExecutionId) zombiecounter += 1 @@ -791,3 +875,73 @@ func zombiecheck(ctx context.Context, workerTimeout int) error { return nil } + +type ExecutionRequest struct { + ExecutionId string `json:"execution_id"` + Authorization string `json:"authorization"` + HTTPProxy string `json:"http_proxy"` + HTTPSProxy string `json:"https_proxy"` + BaseUrl string `json:"base_url"` + EnvironmentName string `json:"environment_name"` + Timezone string `json:"timezone"` + Cleanup string `json:"cleanup"` + ShufflePassProxyToApp string `json:"shuffle_pass_proxy_to_app"` +} + +func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest) error { + parsedRequest := ExecutionRequest{ + ExecutionId: workflowExecution.ExecutionId, + Authorization: workflowExecution.Authorization, + BaseUrl: os.Getenv("BASE_URL"), + EnvironmentName: os.Getenv("ENVIRONMENT_NAME"), + Timezone: os.Getenv("TZ"), + Cleanup: os.Getenv("CLEANUP"), + HTTPProxy: os.Getenv("HTTP_PROXY"), + HTTPSProxy: os.Getenv("HTTPS_PROXY"), + ShufflePassProxyToApp: os.Getenv("SHUFFLE_PASS_APP_PROXY"), + } + + parsedBaseurl := baseUrl + if strings.Contains(baseUrl, ":") { + baseUrlSplit := strings.Split(baseUrl, ":") + if len(baseUrlSplit) >= 3 { + parsedBaseurl = strings.Join(baseUrlSplit[0:2], ":") + //parsedRequest.BaseUrl = fmt.Sprintf("%s:33333", parsedBaseurl) + } + } + + data, err := json.Marshal(parsedRequest) + if err != nil { + log.Printf("[ERROR] Failed marshalling worker request: %s", err) + return err + } + + streamUrl := fmt.Sprintf("%s:33333/api/v1/execute", parsedBaseurl) + req, err := http.NewRequest( + "POST", + streamUrl, + bytes.NewBuffer([]byte(data)), + ) + + client := &http.Client{} + if err != nil { + log.Printf("[ERROR] Failed creating finishing request: %s", err) + return err + } + + newresp, err := client.Do(req) + if err != nil { + log.Printf("[ERROR] Error running finishing request: %s", err) + return err + } + + body, err := ioutil.ReadAll(newresp.Body) + if err != nil { + log.Printf("[ERROR] Failed reading body: %s", err) + return err + } else { + log.Printf("[INFO] NEWRESP (from backend): %s", string(body)) + } + + return nil +} diff --git a/functions/onprem/worker/build.sh b/functions/onprem/worker/build.sh index 1ea9af93..11350ed7 100644 --- a/functions/onprem/worker/build.sh +++ b/functions/onprem/worker/build.sh @@ -1,5 +1,5 @@ NAME=shuffle-worker -VERSION=0.9.28 +VERSION=0.9.29 echo "Running docker build with $NAME:$VERSION" #CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o worker.bin . diff --git a/functions/onprem/worker/go.sum b/functions/onprem/worker/go.sum index b4a04754..21cea8f2 100644 --- a/functions/onprem/worker/go.sum +++ b/functions/onprem/worker/go.sum @@ -571,6 +571,8 @@ github.com/satori/go.uuid v1.2.0/go.mod h1:dA0hQrYB0VpLJoorglMZABFdXlWrHn1NEOzdh github.com/seccomp/libseccomp-golang v0.9.1/go.mod h1:GbW5+tmTXfcxTToHLXlScSlAvWlF4P2Ca7zGrPiEpWo= github.com/shuffle/shuffle-shared v0.1.19 h1:bZmwdC3gKPFxtKoGEjHGY8C96+xB+cBdHaz6glP8SwA= github.com/shuffle/shuffle-shared v0.1.19/go.mod h1:0QrK51T12CpCj/be8hXduj/RtDnoeaZ3rfogELZE2IU= +github.com/shuffle/shuffle-shared v0.1.20 h1:Wz3DZtFtsd/F3rQuZ4KpIURM2HT1jPwEF3eJg8MX9TE= +github.com/shuffle/shuffle-shared v0.1.20/go.mod h1:0QrK51T12CpCj/be8hXduj/RtDnoeaZ3rfogELZE2IU= github.com/shurcooL/sanitized_anchor_name v1.0.0/go.mod h1:1NzhyTcUVG4SuEtjjoZeVRXNmyL/1OwPU0+IJeTBvfc= github.com/sirupsen/logrus v1.0.4-0.20170822132746-89742aefa4b2/go.mod h1:pMByvHTf9Beacp5x1UXfOR9xyW/9antXMhjMPG0dEzc= github.com/sirupsen/logrus v1.0.6/go.mod h1:pMByvHTf9Beacp5x1UXfOR9xyW/9antXMhjMPG0dEzc= diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 2f54a1b8..f02be254 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -31,6 +31,7 @@ import ( "github.com/gorilla/mux" "github.com/patrickmn/go-cache" + "github.com/satori/go.uuid" ) // This is getting out of hand :) @@ -59,6 +60,8 @@ var startAction string var results []shuffle.ActionResult var allLogs map[string]string +var executionRunning bool + // 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 %#v. Result amount: %d. ResultsSent: %d, Send result: %#v", workflowExecution.Status, reason, len(workflowExecution.Results), requestsSent, handleResultSend) @@ -168,8 +171,25 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason //Finished shutdown (after %d seconds). ", sleepDuration) // Allows everything to finish in subprocesses (apps) - time.Sleep(time.Duration(sleepDuration) * time.Second) - os.Exit(3) + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" { + time.Sleep(time.Duration(sleepDuration) * time.Second) + os.Exit(3) + } else { + log.Printf("[DEBUG] Sending result and resetting values (K8s & Swarm).") + 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{} + executionRunning = false + } + //cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId) } // Deploys the internal worker whenever something happens @@ -186,8 +206,12 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] Type: "json-file", Config: map[string]string{}, }, - Resources: container.Resources{}, - NetworkMode: container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId)), + Resources: container.Resources{}, + } + + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" { + hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId)) + log.Printf("Environments: %#v", env) } // Removing because log extraction should happen first @@ -223,8 +247,6 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } else { log.Printf("[WARNING] No mounted folders") } - // hostConfig.Binds = volumeBinds - //} config := &container.Config{ Image: image, @@ -241,11 +263,29 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] ) if err != nil { + //log.Printf("[ERROR] Failed creating container: %s", err) if !strings.Contains(err.Error(), "Conflict. The container name") { - log.Printf("[ERROR] Container CREATE error: %s", err) - } + log.Printf("[ERROR] Container CREATE error (1): %s", err) - return err + return err + } else { + parsedUuid := uuid.NewV4() + identifier = fmt.Sprintf("%s-%s", identifier, parsedUuid) + log.Printf("[INFO] 2 - Identifier: %s", identifier) + cont, err = cli.ContainerCreate( + context.Background(), + config, + hostConfig, + nil, + nil, + identifier, + ) + + if err != nil { + log.Printf("[ERROR] Container create error (2): %s", err) + return err + } + } } err = cli.ContainerStart(ctx, cont.ID, types.ContainerStartOptions{}) @@ -294,12 +334,14 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } if exit { - log.Printf("ERROR IN CONTAINER DEPLOYMENT - ITS EXITED!") + log.Printf("[DEBUG] ERROR IN CONTAINER DEPLOYMENT - ITS EXITED!") return errors.New(fmt.Sprintf(`{"success": false, "reason": "Container %s exited prematurely.","debug": "docker logs -f %s"}`, cont.ID, cont.ID)) } } } + log.Printf("[DEBUG] Deployed container ID %s", cont.ID) + /* //log.Printf("%#v", stats.Config.Status) //ContainerJSONtoConfig(cj dockType.ContainerJSON) ContainerConfig { @@ -770,10 +812,10 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { err = runUserInput(topClient, action, workflowExecution.Workflow.ID, workflowExecution.ExecutionId, workflowExecution.Authorization, string(triggerData)) if err != nil { - log.Printf("Failed launching backend magic: %s", err) + log.Printf("[ERROR] Failed launching backend magic: %s", err) os.Exit(3) } else { - log.Printf("Launched user input node succesfully!") + log.Printf("[INFO] Launched user input node succesfully!") os.Exit(3) } @@ -880,7 +922,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if err != nil || stats.ContainerJSONBase.State.Status != "running" { // REMOVE if err == nil { - log.Printf("Status: %s, should kill: %s", stats.ContainerJSONBase.State.Status, identifier) + log.Printf("[DEBUG] Status: %s, should kill: %s", stats.ContainerJSONBase.State.Status, identifier) err = removeContainer(identifier) if err != nil { log.Printf("Error killing container: %s", err) @@ -906,19 +948,19 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { actionData, err := json.Marshal(action) if err != nil { - log.Printf("Failed unmarshalling action: %s", err) + log.Printf("[WARNING] Failed unmarshalling action: %s", err) continue } if action.AppID == "0ca8887e-b4af-4e3e-887c-87e9d3bc3d3e" { - log.Printf("\nShould run filter: %#v\n\n", action) + log.Printf("[DEBUG] Should run filter: %#v\n\n", action) runFilter(workflowExecution, action) continue } executionData, err := json.Marshal(workflowExecution) if err != nil { - log.Printf("Failed marshalling executiondata: %s", err) + log.Printf("[ERROR] Failed marshalling executiondata: %s", err) executionData = []byte("") } @@ -942,7 +984,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { } // Fixes issue: - // standard_init_linux.go:185: exec user process caused "argument list too long" + // standard_go init_linux.go:185: exec user process caused "argument list too long" // https://devblogs.microsoft.com/oldnewthing/20100203-00/?p=15083 // FIXME: Ensure to NEVER do this anymore @@ -976,6 +1018,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (2)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } err := downloadDockerImageBackend(topClient, image) @@ -988,6 +1031,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (41)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } } else { executed = true @@ -1001,6 +1045,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (3)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } //log.Printf("[WARNING] Failed CLEANUP execution. Downloading image %s remotely.", image) @@ -1012,6 +1057,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } buildBuf := new(strings.Builder) @@ -1020,11 +1066,13 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } 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) + return } log.Printf("[INFO] Successfully downloaded %s", image) @@ -1037,6 +1085,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (7)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } if strings.Contains(err.Error(), "No such image") { @@ -1044,6 +1093,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } } } @@ -1056,6 +1106,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (9)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } // Trying to replace with lowercase to deploy again. This seems to work with Dockerhub well. @@ -1066,6 +1117,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (10)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } log.Printf("[DEBUG] Failed deploy. Downloading image %s", image) @@ -1079,6 +1131,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (40)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } } else { executed = true @@ -1092,6 +1145,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (11)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } log.Printf("[WARNING] Failed deploying image THREE TIMES. Attempting to download %s as last resort from backend and dockerhub.", image) @@ -1101,6 +1155,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } buildBuf := new(strings.Builder) @@ -1109,11 +1164,13 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } 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("Error deploying container: %s", buildBuf.String()), true) + return } log.Printf("[INFO] Successfully downloaded %s", image) @@ -1126,6 +1183,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (15)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) + return } if strings.Contains(err.Error(), "No such image") { @@ -1133,6 +1191,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { 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) + return } } } @@ -1174,6 +1233,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { validateFinished(workflowExecution) log.Printf("[DEBUG] Shutting down (17)") shutdown(workflowExecution, "", "", true) + return } } @@ -1597,7 +1657,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { } if workflowExecution.Status == "FINISHED" { - log.Printf("Workflowexecution is already FINISHED. No further action can be taken") + 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))) return @@ -1734,7 +1794,7 @@ func getWorkflowExecution(ctx context.Context, id string) (*shuffle.WorkflowExec } func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { - if workflowExecution.ExecutionSource == "default" { + if workflowExecution.ExecutionSource == "default" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" { log.Printf("[INFO] Not sending backend info since source is default") return } @@ -1774,7 +1834,7 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) { //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) + log.Printf("[FINISHED] Should send full result to %s", baseUrl) //data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization) shutdownData, err := json.Marshal(workflowExecution) @@ -1858,7 +1918,6 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo } else { log.Printf("[DEBUG] NOT shutting down with dbSave (%s)", workflowExecution.ExecutionSource) } - } return nil @@ -1900,15 +1959,27 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { // container being launched and port being assigned to webserver listener, err := getAvailablePort() if err != nil { - log.Printf("Failed to created listener: %s", err) - log.Printf("[DEBUG] Shutting down (26)") - shutdown(workflowExecution, "", "", true) + log.Printf("[ERROR] Failed to create init listener: %s", err) + return listener } - port := listener.Addr().(*net.TCPAddr).Port - log.Printf("\n\nStarting webserver on port %d with hostname: %s\n\n", port, hostname) log.Printf("OLD HOSTNAME: %s", appCallbackUrl) - appCallbackUrl = fmt.Sprintf("http://%s:%d", hostname, port) + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + log.Printf("\n\nStarting webserver on port 33333 with hostname: %s\n\n", hostname) + appCallbackUrl = fmt.Sprintf("http://%s:33333", hostname) + listener, err = net.Listen("tcp", ":33333") + if err != nil { + log.Printf("[ERROR] Failed to assign port to 33333") + return nil + } + + return listener + } else { + port := listener.Addr().(*net.TCPAddr).Port + + log.Printf("\n\nStarting webserver on port %d with hostname: %s\n\n", port, hostname) + appCallbackUrl = fmt.Sprintf("http://%s:%d", hostname, port) + } log.Printf("NEW HOSTNAME: %s", appCallbackUrl) return listener @@ -1918,6 +1989,13 @@ func runWebserver(listener net.Listener) { r := mux.NewRouter() r.HandleFunc("/api/v1/streams", handleWorkflowQueue).Methods("POST") r.HandleFunc("/api/v1/streams/results", handleGetStreamResults).Methods("POST", "OPTIONS") + + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + requestCache = cache.New(5*time.Minute, 10*time.Minute) + log.Printf("[DEBUG] Running webserver config for SWARM and K8s") + r.HandleFunc("/api/v1/execute", handleRunExecution).Methods("POST", "OPTIONS") + } + http.Handle("/", r) //log.Fatal(http.ListenAndServe(port, nil)) @@ -2028,6 +2106,22 @@ func main() { } log.Printf("[INFO] Running with timezone %s", timezone) + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + workflowExecution := shuffle.WorkflowExecution{} + listener := webserverSetup(workflowExecution) + //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) + //} + //go func() { + // time.Sleep(time.Duration(1)) + // handleExecutionResult(workflowExecution) + //}() + + runWebserver(listener) + } //imageName := fmt.Sprintf("%s/%s:shuffle_openapi_1.0.0", registryName, baseimagename) //downloadDockerImageBackend(client, imageName) @@ -2042,7 +2136,6 @@ func main() { log.Printf("[WARNING] Running test environment for worker by executing workflow %s", testing) authorization, executionId = runTestExecution(client, testing, shuffle_apikey) - //os.Exit(3) } else { authorization = os.Getenv("AUTHORIZATION") executionId = os.Getenv("EXECUTIONID") @@ -2077,6 +2170,7 @@ func main() { log.Printf("[DEBUG] Shutting down (29)") shutdown(workflowExecution, "", "", true) } + topClient = client firstRequest := true @@ -2174,12 +2268,29 @@ func main() { //wg.Add(1) //wg.Wait() } else { - log.Printf("\n\n[INFO] Running NON-OPTIMIZED execution for type %s with %d environments. This only happens when ran manually. Status: %s\n\n", workflowExecution.ExecutionSource, len(environments), workflowExecution.Status) + log.Printf("\n\n[INFO] Running NON-OPTIMIZED execution for type %s with %d environment(s). This only happens when ran manually OR when running with subflows. Status: %s\n\n", workflowExecution.ExecutionSource, len(environments), workflowExecution.Status) //err := executionInit(workflowExecution) //if err != nil { // log.Printf("[INFO] Workflow setup failed: %s", workflowExecution.ExecutionId, err) // shutdown(workflowExecution, "", "", true) //} + + // Trying to make worker into microservice~ :) + if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" { + listener := webserverSetup(workflowExecution) + 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) + } + go func() { + time.Sleep(time.Duration(1)) + handleExecutionResult(workflowExecution) + }() + + runWebserver(listener) + } } } @@ -2206,3 +2317,161 @@ func main() { time.Sleep(time.Duration(sleepTime) * time.Second) } } + +func handleRunExecution(resp http.ResponseWriter, request *http.Request) { + if executionRunning { + log.Println("[WARNING] An execution is already running on this worker") + resp.WriteHeader(500) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "An execution is already running"}`))) + return + } + + executionRunning = true + body, err := ioutil.ReadAll(request.Body) + if err != nil { + executionRunning = false + log.Println("[WARNING] Failed reading body for stream result queue") + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + type ExecutionRequest struct { + ExecutionId string `json:"execution_id"` + Authorization string `json:"authorization"` + HTTPProxy string `json:"http_proxy"` + HTTPSProxy string `json:"https_proxy"` + ShufflePassProxyToApp string `json:"shuffle_pass_proxy_to_app` + BaseUrl string `json:"base_url"` + EnvironmentName string `json:"environment_name"` + Timezone string `json:"timezone"` + Cleanup string `json:"cleanup"` + } + + var execRequest ExecutionRequest + err = json.Unmarshal(body, &execRequest) + if err != nil { + executionRunning = false + log.Printf("[WARNING] Failed shuffle.WorkflowExecution unmarshaling: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + //if strings.ToLower(os.Getenv("SHUFFLE_PASS_APP_PROXY")) == "true" { + if len(execRequest.HTTPProxy) > 0 { + log.Printf("[DEBUG] Sending proxy info to child process") + os.Setenv("SHUFFLE_PASS_APP_PROXY", execRequest.ShufflePassProxyToApp) + } + if len(execRequest.HTTPProxy) > 0 { + log.Printf("[DEBUG] Running with default HTTP proxy %s", execRequest.HTTPProxy) + os.Setenv("HTTP_PROXY", execRequest.HTTPProxy) + } + if len(execRequest.HTTPSProxy) > 0 { + log.Printf("[DEBUG] Running with default HTTPS proxy %s", execRequest.HTTPSProxy) + os.Setenv("HTTPS_PROXY", execRequest.HTTPSProxy) + } + if len(execRequest.EnvironmentName) > 0 { + os.Setenv("ENVIRONMENT_NAME", execRequest.EnvironmentName) + environment = execRequest.EnvironmentName + } + if len(execRequest.Timezone) > 0 { + os.Setenv("TZ", execRequest.Timezone) + timezone = execRequest.Timezone + } + if len(execRequest.Cleanup) > 0 { + os.Setenv("CLEANUP", execRequest.Cleanup) + cleanupEnv = execRequest.Cleanup + } + if len(execRequest.BaseUrl) > 0 { + os.Setenv("BASE_URL", execRequest.BaseUrl) + baseUrl = execRequest.BaseUrl + } + + topClient = &http.Client{} + var workflowExecution shuffle.WorkflowExecution + data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, execRequest.ExecutionId, execRequest.Authorization) + streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl) + + req, err := http.NewRequest( + "POST", + streamResultUrl, + bytes.NewBuffer([]byte(data)), + ) + + newresp, err := topClient.Do(req) + if err != nil { + executionRunning = false + log.Printf("[ERROR] Failed making request: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + body, err = ioutil.ReadAll(newresp.Body) + if err != nil { + executionRunning = false + log.Printf("[ERROR] Failed reading body: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + if newresp.StatusCode != 200 { + executionRunning = false + log.Printf("[ERROR] Bad statuscode: %d, %s", newresp.StatusCode, string(body)) + + if strings.Contains(string(body), "Workflowexecution is already finished") { + log.Printf("[DEBUG] Shutting down (19)") + //shutdown(workflowExecution, "", "", true) + } + + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad statuscode: %d"}`, newresp.StatusCode))) + return + } + + err = json.Unmarshal(body, &workflowExecution) + if err != nil { + executionRunning = false + log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + return + } + + if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" { + 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))) + return + } + + 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 + log.Printf("[WARNING] Exiting as worker execution has status %s!", workflowExecution.Status) + log.Printf("[DEBUG] Shutting down (21)") + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad status %s"}`, workflowExecution.Status))) + return + } + + log.Printf("[DEBUG] Starting execution :O") + resp.WriteHeader(200) + resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) + + cacheKey := fmt.Sprintf("workflowexecution-%s", workflowExecution.ExecutionId) + requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration) + + 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) + } + + handleExecutionResult(workflowExecution) +}