diff --git a/backend/go-app/go.mod b/backend/go-app/go.mod
index c93d1f5f..cbbc0ea3 100644
--- a/backend/go-app/go.mod
+++ b/backend/go-app/go.mod
@@ -18,7 +18,7 @@ require (
github.com/gorilla/mux v1.8.0
github.com/h2non/filetype v1.1.3
github.com/satori/go.uuid v1.2.0
- github.com/shuffle/shuffle-shared v0.5.14
+ github.com/shuffle/shuffle-shared v0.5.29
golang.org/x/crypto v0.14.0
google.golang.org/api v0.125.0
google.golang.org/grpc v1.55.0
diff --git a/backend/go-app/go.sum b/backend/go-app/go.sum
index 4812fee4..bea5b3c6 100644
--- a/backend/go-app/go.sum
+++ b/backend/go-app/go.sum
@@ -420,6 +420,8 @@ github.com/shuffle/shuffle-shared v0.5.11 h1:Eqbs9o8E49QAL5/6aV6BfFtWSjLIvgET7AL
github.com/shuffle/shuffle-shared v0.5.11/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/shuffle/shuffle-shared v0.5.14 h1:d14u1e4k+qKgnf4Insq4x2S+0MMKlDqdyTTyVP3puRA=
github.com/shuffle/shuffle-shared v0.5.14/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
+github.com/shuffle/shuffle-shared v0.5.29 h1:n4vThl7v3mFVXbrIW71XREFdmZZo7mOBAWxnsdiNjDk=
+github.com/shuffle/shuffle-shared v0.5.29/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/shurcooL/sanitized_anchor_name v1.0.0/go.mod h1:1NzhyTcUVG4SuEtjjoZeVRXNmyL/1OwPU0+IJeTBvfc=
github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0=
github.com/sirupsen/logrus v1.8.1 h1:dJKuHgqk1NNQlqoA6BTlM1Wf9DOH3NBjQyu0h9+AZZE=
diff --git a/frontend/src/components/NewHeader.jsx b/frontend/src/components/NewHeader.jsx
index 42d79ebd..77cabb46 100644
--- a/frontend/src/components/NewHeader.jsx
+++ b/frontend/src/components/NewHeader.jsx
@@ -610,6 +610,11 @@ const Header = (props) => {
>
Logout
+
+
+
+ Version: 1.3.1
+
);
diff --git a/frontend/src/components/ParsedAction.jsx b/frontend/src/components/ParsedAction.jsx
index 558fb983..e0abbe42 100755
--- a/frontend/src/components/ParsedAction.jsx
+++ b/frontend/src/components/ParsedAction.jsx
@@ -799,7 +799,6 @@ const ParsedAction = (props) => {
selectedActionParameters[count].value = splitparsed[0]
selectedAction.parameters[count].value = splitparsed[0]
- //changeActionParameter({target: {value: splitparsed[1]}},
selectedActionParameters[1].value = splitparsed[1]
selectedAction.parameters[1].value = splitparsed[1]
forceUpdate = true
diff --git a/frontend/src/components/RuntimeDebugger.jsx b/frontend/src/components/RuntimeDebugger.jsx
index 1b4f431f..5810cbab 100644
--- a/frontend/src/components/RuntimeDebugger.jsx
+++ b/frontend/src/components/RuntimeDebugger.jsx
@@ -478,7 +478,7 @@ const RuntimeDebugger = (props) => {
return (
-
Workflow result: {errorReason}
{params.row.result !== null && params.row.result !== undefined && params.row.result !== "" ?
diff --git a/frontend/src/views/AngularWorkflow.jsx b/frontend/src/views/AngularWorkflow.jsx
index 0b710867..e9e28582 100755
--- a/frontend/src/views/AngularWorkflow.jsx
+++ b/frontend/src/views/AngularWorkflow.jsx
@@ -71,6 +71,7 @@ import {
import {
Folder as FolderIcon,
+ Insights as InsightsIcon,
LibraryBooks as LibraryBooksIcon,
OpenInNew as OpenInNewIcon,
Undo as UndoIcon,
@@ -15193,6 +15194,27 @@ const AngularWorkflow = (defaultprops) => {
) : null}
+ {isCloud ? (
+
+
+
+
+
+ ) : null}
{executionData.workflow !== undefined && executionData.workflow !== null && executionData.workflow.actions !== undefined && executionData.workflow.actions !== null && executionData.workflow.actions.length > 0 && executionData.workflow.actions[0].environment !== "Cloud" ?
diff --git a/frontend/src/views/Workflows.jsx b/frontend/src/views/Workflows.jsx
index 981bc7a0..73f02cc2 100755
--- a/frontend/src/views/Workflows.jsx
+++ b/frontend/src/views/Workflows.jsx
@@ -1344,7 +1344,7 @@ const Workflows = (props) => {
}, i * 200);
}
- toast(`exporting and keeping original for all ${allWorkflows.length} workflows`);
+ toast(`Exporting and keeping original for all ${allWorkflows.length} workflows`);
};
const deduplicateIds = (data, skip_sanitize) => {
diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod
index fef487a1..4d2d9450 100644
--- a/functions/onprem/orborus/go.mod
+++ b/functions/onprem/orborus/go.mod
@@ -9,7 +9,7 @@ require (
github.com/mackerelio/go-osstat v0.2.3
github.com/satori/go.uuid v1.2.0
github.com/shirou/gopsutil v3.21.11+incompatible
- github.com/shuffle/shuffle-shared v0.4.96
+ github.com/shuffle/shuffle-shared v0.5.29
k8s.io/api v0.28.1
k8s.io/apimachinery v0.28.1
k8s.io/client-go v0.28.1
diff --git a/functions/onprem/orborus/go.sum b/functions/onprem/orborus/go.sum
index c9eea384..fa8b721c 100644
--- a/functions/onprem/orborus/go.sum
+++ b/functions/onprem/orborus/go.sum
@@ -278,6 +278,8 @@ github.com/shuffle/shuffle-shared v0.4.95 h1:xr92/03/uQeJiDme9S8/vgF1KWyQgJ1KQXV
github.com/shuffle/shuffle-shared v0.4.95/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/shuffle/shuffle-shared v0.4.96 h1:iaIB/HP9eKpw9DMMJZhSLDbKdHJt075kFYLHg9AaiiM=
github.com/shuffle/shuffle-shared v0.4.96/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
+github.com/shuffle/shuffle-shared v0.5.29 h1:n4vThl7v3mFVXbrIW71XREFdmZZo7mOBAWxnsdiNjDk=
+github.com/shuffle/shuffle-shared v0.5.29/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0=
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
diff --git a/functions/onprem/worker/go.mod b/functions/onprem/worker/go.mod
index e4cf2262..9e0f8e1e 100644
--- a/functions/onprem/worker/go.mod
+++ b/functions/onprem/worker/go.mod
@@ -11,7 +11,7 @@ require (
github.com/gorilla/mux v1.8.0
github.com/patrickmn/go-cache v2.1.0+incompatible
github.com/satori/go.uuid v1.2.0
- github.com/shuffle/shuffle-shared v0.4.57
+ github.com/shuffle/shuffle-shared v0.5.29
k8s.io/api v0.28.3
k8s.io/apimachinery v0.28.3
k8s.io/client-go v0.28.3
diff --git a/functions/onprem/worker/go.sum b/functions/onprem/worker/go.sum
index f26caf60..5afdce76 100644
--- a/functions/onprem/worker/go.sum
+++ b/functions/onprem/worker/go.sum
@@ -284,6 +284,8 @@ github.com/shuffle/shuffle-shared v0.4.50 h1:fJLfhWIJ5mYap4JwHnD/B5aaLyIULwylFSl
github.com/shuffle/shuffle-shared v0.4.50/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/shuffle/shuffle-shared v0.4.57 h1:o+mMPRY4ourkE3R0qdi80jg6RlCtvAJ/VVrPk4y75Hk=
github.com/shuffle/shuffle-shared v0.4.57/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
+github.com/shuffle/shuffle-shared v0.5.29 h1:n4vThl7v3mFVXbrIW71XREFdmZZo7mOBAWxnsdiNjDk=
+github.com/shuffle/shuffle-shared v0.5.29/go.mod h1:X613gbo0dT3fnYvXDRwjQZyLC+T49T2nSQOrCV5QMlI=
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e h1:MRM5ITcdelLK2j1vwZ3Je0FKVCfqOLp5zO6trqMLYs0=
github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e/go.mod h1:XV66xRDqSt+GTGFMVlhk3ULuV0y9ZmzeVGR4mloJI3M=
github.com/spf13/pflag v1.0.5 h1:iy+VFUOCP1a+8yFto/drg2CJ5u0yRoB7fZw3DKv/JXA=
diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go
index fb913d16..0c8fd3a5 100755
--- a/functions/onprem/worker/worker.go
+++ b/functions/onprem/worker/worker.go
@@ -3,7 +3,7 @@ package main
import (
"github.com/shuffle/shuffle-shared"
- //"bufio"
+
"bytes"
"context"
"encoding/json"
@@ -14,29 +14,23 @@ import (
"log"
"net"
"net/http"
+ "net/http/pprof"
"net/url"
"os"
+ "strconv"
"strings"
"time"
"github.com/docker/docker/api/types"
+ "github.com/docker/docker/api/types/filters"
"github.com/docker/docker/api/types/container"
- //"github.com/docker/docker/api/types/filters"
"github.com/docker/docker/api/types/mount"
dockerclient "github.com/docker/docker/client"
- //"github.com/go-git/go-billy/v5/memfs"
-
- //newdockerclient "github.com/fsouza/go-dockerclient"
- //"github.com/satori/go.uuid"
+ // This is for automatic removal of certain code :)
"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"
-
//k8s deps
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
@@ -45,24 +39,24 @@ import (
"k8s.io/client-go/tools/clientcmd"
"k8s.io/client-go/util/homedir"
"path/filepath"
- // "k8s.io/client-go/util/retry"
)
// This is getting out of hand :)
-var environment = os.Getenv("ENVIRONMENT_NAME")
+var timezone = os.Getenv("TZ")
var baseUrl = os.Getenv("BASE_URL")
var appCallbackUrl = os.Getenv("BASE_URL")
+var isKubernetes = os.Getenv("IS_KUBERNETES")
+var environment = os.Getenv("ENVIRONMENT_NAME")
+var logsDisabled = os.Getenv("SHUFFLE_LOGS_DISABLED")
var cleanupEnv = strings.ToLower(os.Getenv("CLEANUP"))
-var dockerApiVersion = strings.ToLower(os.Getenv("DOCKER_API_VERSION"))
var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
-var timezone = os.Getenv("TZ")
+var dockerApiVersion = strings.ToLower(os.Getenv("DOCKER_API_VERSION"))
var baseimagename = "frikky/shuffle"
// var baseimagename = "registry.hub.docker.com/frikky/shuffle"
var registryName = "registry.hub.docker.com"
var sleepTime = 2
-var requestCache *cache.Cache
var topClient *http.Client
var data string
var requestsSent = 0
@@ -84,11 +78,21 @@ var startAction string
//var allLogs map[string]string
//var containerIds []string
var downloadedImages []string
+type ImageDownloadBody struct {
+ Image string `json:"image"`
+}
+
+type ImageRequest struct {
+ Image string `json:"image"`
+}
+
+var finishedExecutions []string
+
// Images to be autodeployed in the latest version of Shuffle.
var autoDeploy = map[string]string{
- "http:1.3.0": "frikky/shuffle:http_1.3.0",
"http:1.4.0": "frikky/shuffle:http_1.4.0",
+ "http:1.3.0": "frikky/shuffle:http_1.3.0",
"shuffle-tools:1.2.0": "frikky/shuffle:shuffle-tools_1.2.0",
"shuffle-subflow:1.0.0": "frikky/shuffle:shuffle-subflow_1.0.0",
"shuffle-subflow:1.1.0": "frikky/shuffle:shuffle-subflow_1.1.0",
@@ -108,6 +112,166 @@ type UserInputSubflow struct {
CancelUrl string `json:"cancel_url"`
}
+// Not using shuffle.SetWorkflowExecution as we only want to use cache in reality
+func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.WorkflowExecution, dbSave bool) error {
+ if len(workflowExecution.ExecutionId) == 0 {
+ log.Printf("[DEBUG] Workflowexecution executionId can't be empty.")
+ return errors.New("ExecutionId can't be empty.")
+ }
+
+ //log.Printf("[DEBUG][%s] Setting with %d results (pre)", workflowExecution.ExecutionId, len(workflowExecution.Results))
+ workflowExecution = shuffle.Fixexecution(ctx, workflowExecution)
+ cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId)
+
+ execData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed marshalling execution during set: %s", err)
+ return err
+ }
+
+ err = shuffle.SetCache(ctx, cacheKey, execData, 30)
+ if err != nil {
+ log.Printf("[ERROR][%s] Failed adding to cache during setexecution", workflowExecution)
+ return err
+ }
+
+
+ handleExecutionResult(workflowExecution)
+ validated := shuffle.ValidateFinished(ctx, -1, workflowExecution)
+ if validated {
+ shutdownData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed marshalling shutdowndata during set: %s", err)
+ }
+
+ log.Printf("[DEBUG][%s] Sending result (set)", workflowExecution.ExecutionId)
+ sendResult(workflowExecution, shutdownData)
+ return nil
+ }
+
+ // FIXME: Should this shutdown OR send the result?
+ // The worker may not be running the backend hmm
+ if dbSave {
+ if workflowExecution.ExecutionSource == "default" {
+ log.Printf("[DEBUG][%s] Shutting down (25)", workflowExecution.ExecutionId)
+ shutdown(workflowExecution, "", "", true)
+ //return
+ } else {
+ log.Printf("[DEBUG][%s] NOT shutting down with dbSave (%s). Instead sending result to backend and start polling until subflow is updated", workflowExecution.ExecutionId, workflowExecution.ExecutionSource)
+
+ shutdownData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed marshalling shutdowndata during dbSave handler: %s", err)
+ }
+
+ sendResult(workflowExecution, shutdownData)
+
+ // Poll for 1 minute max if there is a "wait for results" subflow
+ subflowId := ""
+ for _, result := range workflowExecution.Results {
+ if result.Status == "WAITING" {
+ //log.Printf("[DEBUG][%s] Found waiting result", workflowExecution.ExecutionId)
+ subflowId = result.Action.ID
+ }
+ }
+
+ if len(subflowId) == 0 {
+ log.Printf("[DEBUG][%s] No waiting result found. Not polling", workflowExecution.ExecutionId)
+
+ for _, action := range workflowExecution.Workflow.Actions {
+ if action.AppName == "User Input" || action.AppName == "Shuffle Workflow" || action.AppName == "shuffle-subflow" {
+ workflowExecution.Workflow.Triggers = append(workflowExecution.Workflow.Triggers, shuffle.Trigger{
+ AppName: action.AppName,
+ Parameters: action.Parameters,
+ ID: action.ID,
+ })
+ }
+ }
+
+
+ for _, trigger := range workflowExecution.Workflow.Triggers {
+ //log.Printf("[DEBUG] Found trigger %s", trigger.AppName)
+ if trigger.AppName != "User Input" && trigger.AppName != "Shuffle Workflow" && trigger.AppName != "shuffle-subflow" {
+ continue
+ }
+
+
+ // check if it has wait for results in params
+ wait := false
+ for _, param := range trigger.Parameters {
+ //log.Printf("[DEBUG] Found param %s with value %s", param.Name, param.Value)
+ if param.Name == "check_result" && strings.ToLower(param.Value) == "true" {
+ //log.Printf("[DEBUG][%s] Found check result param!", workflowExecution.ExecutionId)
+ wait = true
+ break
+ }
+ }
+
+ if wait {
+ // Check if it has a result or not
+ found := false
+ for _, result := range workflowExecution.Results {
+ //log.Printf("[DEBUG][%s] Found result %s", workflowExecution.ExecutionId, result.Action.ID)
+ if result.Action.ID == trigger.ID && result.Status != "SUCCESS" && result.Status != "FAILURE" {
+ //log.Printf("[DEBUG][%s] Found subflow result that is not handled. Waiting for results", workflowExecution.ExecutionId)
+
+ subflowId = result.Action.ID
+ found = true
+ break
+ }
+ }
+
+ if !found {
+ log.Printf("[DEBUG][%s] No result found for subflow. Setting subflowId to %s", workflowExecution.ExecutionId, trigger.ID)
+ subflowId = trigger.ID
+ }
+ }
+
+ if len(subflowId) > 0 {
+ break
+ }
+ }
+ }
+
+ if len(subflowId) > 0 {
+ // Under rerun period timeout
+ timeComparison := 120
+ log.Printf("[DEBUG][%s] Starting polling for %d seconds to see if new subflow updates are found on the backend that are not handled. Subflow ID: %s", workflowExecution.ExecutionId, timeComparison, subflowId)
+ timestart := time.Now()
+ streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl)
+ for {
+ err = handleSubflowPoller(ctx, workflowExecution, streamResultUrl, subflowId)
+ if err == nil {
+ log.Printf("[DEBUG] Subflow is finished and we are breaking the thingy")
+
+ if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" && workflowExecution.ExecutionSource != "default" {
+ log.Printf("[DEBUG] Force shutdown of worker due to optimized run with webserver. Expecting reruns to take care of this")
+ os.Exit(0)
+ }
+
+
+ break
+ }
+
+ timepassed := time.Since(timestart)
+ if timepassed.Seconds() > float64(timeComparison) {
+ log.Printf("[DEBUG][%s] Max poll time reached to look for updates. Stopping poll. This poll is here to send personal results back to itself to be handled, then to stop this thread.", workflowExecution.ExecutionId)
+ break
+ }
+
+ // Sleep for 1 second
+ time.Sleep(1 * time.Second)
+ }
+ } else {
+ log.Printf("[DEBUG][%s] No need to poll for results. Not polling", workflowExecution.ExecutionId)
+ }
+ }
+ }
+
+ return nil
+}
+
+
// removes every container except itself (worker)
func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason string, handleResultSend bool) {
log.Printf("[DEBUG][%s] Shutdown (%s) started with reason %#v. Result amount: %d. ResultsSent: %d, Send result: %#v, Parent: %#v", workflowExecution.ExecutionId, workflowExecution.Status, reason, len(workflowExecution.Results), requestsSent, handleResultSend, workflowExecution.ExecutionParent)
@@ -126,6 +290,30 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
time.Sleep(time.Duration(sleepDuration) * time.Second)
}
+ // Might not be necessary because of cleanupEnv hostconfig autoremoval
+ if cleanupEnv == "true" && (os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm") {
+ /*
+ ctx := context.Background()
+ dockercli, err := dockerclient.NewEnvClient()
+ if err == nil {
+ log.Printf("[INFO] Cleaning up %d containers", len(containerIds))
+ removeOptions := types.ContainerRemoveOptions{
+ RemoveVolumes: true,
+ Force: true,
+ }
+
+ for _, containername := range containerIds {
+ log.Printf("[INFO] Should stop and and remove container %s (deprecated)", containername)
+ //dockercli.ContainerStop(ctx, containername, nil)
+ //dockercli.ContainerRemove(ctx, containername, removeOptions)
+ //removeContainers = append(removeContainers, containername)
+ }
+ }
+ */
+ } else {
+
+ }
+
if len(reason) > 0 && len(nodeId) > 0 {
//log.Printf("[INFO] Running abort of workflow because it should be finished")
@@ -138,7 +326,6 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
path += fmt.Sprintf("&env=%s", url.QueryEscape(environment))
}
- //fmt.Printf(url.QueryEscape(query))
abortUrl += path
log.Printf("[DEBUG][%s] Abort URL: %s", workflowExecution.ExecutionId, abortUrl)
@@ -152,19 +339,22 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
log.Printf("[WARNING][%s] Failed building request: %s", workflowExecution.ExecutionId, err)
}
- authorization := os.Getenv("AUTHORIZATION")
- if len(authorization) > 0 {
- req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization))
+ // FIXME: Add an API call to the backend
+ if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
+ authorization := os.Getenv("AUTHORIZATION")
+ if len(authorization) > 0 {
+ req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization))
+ } else {
+ log.Printf("[ERROR][%s] No authorization specified for abort", workflowExecution.ExecutionId)
+ }
} else {
- log.Printf("[ERROR][%s] No authorization specified for abort", workflowExecution.ExecutionId)
+ req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", workflowExecution.Authorization))
}
req.Header.Add("Content-Type", "application/json")
- client := shuffle.GetExternalClient(baseUrl)
-
//log.Printf("[DEBUG][%s] All App Logs: %#v", workflowExecution.ExecutionId, allLogs)
- newresp, err := client.Do(req)
+ newresp, err := topClient.Do(req)
if err != nil {
log.Printf("[WARNING][%s] Failed abort request: %s", workflowExecution.ExecutionId, err)
} else {
@@ -178,63 +368,17 @@ 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)
-}
-
-// }
-
-func isRunningInCluster() bool {
- _, existsHost := os.LookupEnv("KUBERNETES_SERVICE_HOST")
- _, existsPort := os.LookupEnv("KUBERNETES_SERVICE_PORT")
- return existsHost && existsPort
-}
-
-func buildEnvVars(envMap map[string]string) []corev1.EnvVar {
- var envVars []corev1.EnvVar
- for key, value := range envMap {
- envVars = append(envVars, corev1.EnvVar{Name: key, Value: value})
- }
- return envVars
-}
-
-func getKubernetesClient() (*kubernetes.Clientset, error) {
- if isRunningInCluster() {
- config, err := rest.InClusterConfig()
- if err != nil {
- return nil, err
- }
- clientset, err := kubernetes.NewForConfig(config)
- if err != nil {
- return nil, err
- }
- return clientset, nil
+ if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
+ time.Sleep(time.Duration(sleepDuration) * time.Second)
+ os.Exit(3)
} else {
- home := homedir.HomeDir()
- kubeconfigPath := filepath.Join(home, ".kube", "config")
- config, err := clientcmd.BuildConfigFromFlags("", kubeconfigPath)
- if err != nil {
- return nil, err
- }
- clientset, err := kubernetes.NewForConfig(config)
- if err != nil {
- return nil, err
- }
- return clientset, nil
+ log.Printf("[DEBUG][%s] Sending result and resetting values (K8s & Swarm).", workflowExecution.ExecutionId)
}
}
// Deploys the internal worker whenever something happens
func deployApp(cli *dockerclient.Client, image string, identifier string, env []string, workflowExecution shuffle.WorkflowExecution, action shuffle.Action) error {
- // log.Printf("################################### new call to deployApp ###################################")
- // log.Printf("image: %s", image)
- // log.Printf("identifier: %s", identifier)
- // log.Printf("execution: %+v", workflowExecution)
- log.Printf("[DEBUG] Adding SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT"))
- env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
-
- if os.Getenv("IS_KUBERNETES") == "true" {
-
+ if isKubernetes == "true" {
namespace := "shuffle"
localRegistry := os.Getenv("REGISTRY_URL")
@@ -248,8 +392,9 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
clientset, err := getKubernetesClient()
if err != nil {
- fmt.Println("[ERROR]Error getting kubernetes client:", err)
- // os.Exit(1)
+ log.Printf("[ERROR] Failed getting kubernetes: %s [INFO] Setting kubernetes to false to enable running Shuffle with Docker for the next iterations.", err)
+ isKubernetes = "false"
+ return err
}
log.Printf("[DEBUG] Got kubernetes client")
@@ -264,8 +409,6 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
appName := strings.Join(appDetailsSplit[:len(appDetailsSplit)-1], "_")
appVersion := appDetailsSplit[len(appDetailsSplit)-1]
- // log.Printf("APP VERSION IS: %s", appVersion)
-
for _, app := range workflowExecution.Workflow.Actions {
// log.Printf("[DEBUG] App: %s, Version: %s", appName, appVersion)
// log.Printf("[DEBUG] Checking app %s with version %s", app.AppName, app.AppVersion)
@@ -308,115 +451,140 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
createdPod, err := clientset.CoreV1().Pods(namespace).Create(context.Background(), pod, metav1.CreateOptions{})
if err != nil {
- fmt.Fprintf(os.Stderr, "Error creating pod: %v\n", err)
+ fmt.Fprintf(os.Stderr, "Error creating pod: %v", err)
// os.Exit(1)
}
- fmt.Printf("[DEBUG] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace)
- } else {
- // form basic hostConfig
- ctx := context.Background()
-
- if action.AppName == "shuffle-subflow" {
- // Automatic replacement of URL
- for paramIndex, param := range action.Parameters {
- if param.Name != "backend_url" {
- continue
- }
-
- if strings.Contains(param.Value, "shuffle-backend") {
- // Automatic replacement as this is default
- action.Parameters[paramIndex].Value = os.Getenv("BASE_URL")
- log.Printf("[DEBUG][%s] Replaced backend_url with %s", workflowExecution.ExecutionId, os.Getenv("BASE_URL"))
- }
- }
- }
-
- // Max 10% CPU every second
- //CPUShares: 128,
- //CPUQuota: 10000,
- //CPUPeriod: 100000,
- hostConfig := &container.HostConfig{
- LogConfig: container.LogConfig{
- Type: "json-file",
- Config: map[string]string{
- "max-size": "10m",
- },
- },
- Resources: container.Resources{},
- }
-
- hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId))
-
- // Removing because log extraction should happen first
- if cleanupEnv == "true" {
- hostConfig.AutoRemove = true
- }
-
- // FIXME: Add proper foldermounts here
- //log.Printf("\n\nPRE FOLDERMOUNT\n\n")
- //volumeBinds := []string{"/tmp/shuffle-mount:/rules"}
- //volumeBinds := []string{"/tmp/shuffle-mount:/rules"}
- volumeBinds := []string{}
- if len(volumeBinds) > 0 {
- log.Printf("[DEBUG] Setting up binds for container!")
- hostConfig.Binds = volumeBinds
- hostConfig.Mounts = []mount.Mount{}
- for _, bind := range volumeBinds {
- if !strings.Contains(bind, ":") || strings.Contains(bind, "..") || strings.HasPrefix(bind, "~") {
- log.Printf("[WARNING] Bind %s is invalid.", bind)
- continue
- }
-
- log.Printf("[DEBUG] Appending bind %s", bind)
- bindSplit := strings.Split(bind, ":")
- sourceFolder := bindSplit[0]
- destinationFolder := bindSplit[0]
- hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{
- Type: mount.TypeBind,
- Source: sourceFolder,
- Target: destinationFolder,
- })
- }
- } else {
- //log.Printf("[WARNING] Not mounting folders")
- }
-
- config := &container.Config{
- Image: image,
- Env: env,
- }
-
- // Checking as late as possible, just in case.
- newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
- _, err := shuffle.GetCache(ctx, newExecId)
- if err == nil {
- log.Printf("\n\n[DEBUG] Result for %s already found - returning\n\n", newExecId)
- return nil
- }
-
- cacheData := []byte("1")
- err = shuffle.SetCache(ctx, newExecId, cacheData, 30)
- if err != nil {
- log.Printf("[WARNING] Failed setting cache for action %s: %s", newExecId, err)
- } else {
- log.Printf("[DEBUG] Adding %s to cache. Name: %s", newExecId, action.Name)
- }
-
- if action.ExecutionDelay > 0 {
- log.Printf("[DEBUG] Running app %s in docker with delay of %d", action.Name, action.ExecutionDelay)
- waitTime := time.Duration(action.ExecutionDelay) * time.Second
-
- time.AfterFunc(waitTime, func() {
- DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
- })
- } else {
- log.Printf("[DEBUG] Running app %s in docker NORMALLY as there is no delay set with identifier %s", action.Name, identifier)
- returnvalue := DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
- log.Printf("[DEBUG] Normal deploy ret: %s", returnvalue)
- return returnvalue
- }
+ log.Printf("[DEBUG] Created pod %q in namespace %q", createdPod.Name, createdPod.Namespace)
return nil
}
+
+ // form basic hostConfig
+ ctx := context.Background()
+
+ // Check action if subflow
+ // Check if url is default (shuffle-backend)
+ // If it doesn't exist, add it
+ if action.AppName == "shuffle-subflow" {
+ // Automatic replacement of URL
+ for paramIndex, param := range action.Parameters {
+ if param.Name != "backend_url" {
+ continue
+ }
+
+ if strings.Contains(param.Value, "shuffle-backend") {
+ // Automatic replacement as this is default
+ if len(os.Getenv("BASE_URL")) > 0 {
+ action.Parameters[paramIndex].Value = os.Getenv("BASE_URL")
+ log.Printf("[DEBUG][%s] Replaced backend_url with base_url %s", workflowExecution.ExecutionId, os.Getenv("BASE_URL"))
+ }
+
+ if len(os.Getenv("SHUFFLE_CLOUDRUN_URL")) > 0 {
+ action.Parameters[paramIndex].Value = os.Getenv("SHUFFLE_CLOUDRUN_URL")
+ log.Printf("[DEBUG][%s] Replaced backend_url with cloudrun %s", workflowExecution.ExecutionId, os.Getenv("SHUFFLE_CLOUDRUN_URL"))
+ }
+ }
+ }
+ }
+
+
+ // Max 10% CPU every second
+ //CPUShares: 128,
+ //CPUQuota: 10000,
+ //CPUPeriod: 100000,
+ hostConfig := &container.HostConfig{
+ LogConfig: container.LogConfig{
+ Type: "json-file",
+ Config: map[string]string{
+ "max-size": "10m",
+ },
+ },
+ Resources: container.Resources{},
+ }
+
+ if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
+ hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId))
+ //log.Printf("Environments: %#v", env)
+ }
+
+ // Removing because log extraction should happen first
+ if cleanupEnv == "true" {
+ hostConfig.AutoRemove = true
+ }
+
+ // Get environment for certificates
+ volumeBinds := []string{}
+ volumeBindString:= os.Getenv("SHUFFLE_VOLUME_BINDS")
+ if len(volumeBindString) > 0 {
+ volumeBindSplit := strings.Split(volumeBindString, ",")
+ for _, volumeBind := range volumeBindSplit {
+ if strings.Contains(volumeBind, ":") {
+ volumeBinds = append(volumeBinds, volumeBind)
+ } else {
+ log.Printf("[ERROR] Volume bind '%s' is invalid.", volumeBind)
+ }
+ }
+ }
+
+ // Add more volume binds if possible
+ if len(volumeBinds) > 0 {
+ log.Printf("[DEBUG] Setting up binds for container. Got %d volume binds.", len(volumeBinds))
+
+ hostConfig.Binds = volumeBinds
+ hostConfig.Mounts = []mount.Mount{}
+ for _, bind := range volumeBinds {
+ if !strings.Contains(bind, ":") || strings.Contains(bind, "..") || strings.HasPrefix(bind, "~") {
+ log.Printf("[ERROR] Volume bind '%s' is invalid. Use absolute paths.", bind)
+ continue
+ }
+
+ log.Printf("[DEBUG] Appending bind %s to app container", bind)
+ bindSplit := strings.Split(bind, ":")
+ sourceFolder := bindSplit[0]
+ destinationFolder := bindSplit[1]
+ hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{
+ Type: mount.TypeBind,
+ Source: sourceFolder,
+ Target: destinationFolder,
+ })
+ }
+ }
+
+ config := &container.Config{
+ Image: image,
+ Env: env,
+ }
+
+
+ // Checking as late as possible, just in case.
+ newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
+ _, err := shuffle.GetCache(ctx, newExecId)
+ if err == nil {
+ log.Printf("[DEBUG][%s] Result for action %s already found - returning", newExecId, action.ID)
+ return nil
+ }
+
+ cacheData := []byte("1")
+ err = shuffle.SetCache(ctx, newExecId, cacheData, 30)
+ if err != nil {
+ //log.Printf("[WARNING][%s] Failed setting cache for action: %s", newExecId, err)
+ } else {
+ //log.Printf("[DEBUG][%s] Adding to cache. Name: %s", workflowExecution.ExecutionId, action.Name)
+ }
+
+ if action.ExecutionDelay > 0 {
+ log.Printf("[DEBUG][%s] Running app '%s' with label '%s' in docker with delay of %d", workflowExecution.ExecutionId, action.AppName, action.Label, action.ExecutionDelay)
+ waitTime := time.Duration(action.ExecutionDelay) * time.Second
+
+ time.AfterFunc(waitTime, func() {
+ DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
+ })
+ } else {
+ log.Printf("[DEBUG][%s] Running app %s in docker NORMALLY as there is no delay set with identifier %s", workflowExecution.ExecutionId, action.Name, identifier)
+ returnvalue := DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
+ //log.Printf("[DEBUG][%s] Normal deploy ret: %s", workflowExecution.ExecutionId, returnvalue)
+ return returnvalue
+ }
+
return nil
}
@@ -429,7 +597,7 @@ func cleanupExecution(clientset *kubernetes.Clientset, workflowExecution shuffle
LabelSelector: labelSelector,
})
if err != nil {
- return fmt.Errorf("[ERROR]failed to list apps with label selector %s: %v", labelSelector, err)
+ return fmt.Errorf("[ERROR] Failed to list apps with label selector %s: %#vv", labelSelector, err)
}
for _, pod := range podList.Items {
@@ -437,14 +605,14 @@ func cleanupExecution(clientset *kubernetes.Clientset, workflowExecution shuffle
if err != nil {
return fmt.Errorf("failed to delete app %s: %v", pod.Name, err)
}
- fmt.Printf("App %s in namespace %s deleted.\n", pod.Name, namespace)
+ log.Printf("App %s in namespace %s deleted.", pod.Name, namespace)
}
podErr := clientset.CoreV1().Pods(namespace).Delete(context.TODO(), workerName, metav1.DeleteOptions{})
if podErr != nil {
return fmt.Errorf("[ERROR] failed to delete the worker %s in namespace %s: %v", workerName, namespace, podErr)
}
- fmt.Printf("[DEBUG] %s in namespace %s deleted.\n", workerName, namespace)
+ log.Printf("[DEBUG] %s in namespace %s deleted.", workerName, namespace)
return nil
}
@@ -458,6 +626,8 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
identifier,
)
+ //log.Printf("[DEBUG] config set: %#v", config)
+
if err != nil {
//log.Printf("[ERROR] Failed creating container: %s", err)
if !strings.Contains(err.Error(), "Conflict. The container name") {
@@ -465,7 +635,7 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
cacheErr := shuffle.DeleteCache(ctx, newExecId)
if cacheErr != nil {
- log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, cacheErr)
+ log.Printf("[ERROR] FAILURE Deleting cache for %s: %s", newExecId, cacheErr)
}
return err
@@ -489,7 +659,7 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
cacheErr := shuffle.DeleteCache(ctx, newExecId)
if cacheErr != nil {
- log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, cacheErr)
+ log.Printf("[ERROR] FAILURE Deleting cache for %s: %s", newExecId, cacheErr)
}
return err
@@ -528,7 +698,7 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
cacheErr := shuffle.DeleteCache(ctx, newExecId)
if cacheErr != nil {
- log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, cacheErr)
+ log.Printf("[ERROR] FAILURE Deleting cache for %s: %s", newExecId, cacheErr)
}
return err
@@ -543,7 +713,7 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
cacheErr := shuffle.DeleteCache(ctx, newExecId)
if cacheErr != nil {
- log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, cacheErr)
+ log.Printf("[ERROR] FAILURE Deleting cache for %s: %s", newExecId, cacheErr)
}
//shutdown(workflowExecution, workflowExecution.Workflow.ID, true)
@@ -551,16 +721,16 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
}
}
- log.Printf("[DEBUG] Container %s was created for %s", cont.ID, identifier)
+ log.Printf("[DEBUG][%s] Container %s was created for %s", workflowExecution.ExecutionId, cont.ID, identifier)
// Waiting to see if it exits.. Stupid, but stable(r)
if workflowExecution.ExecutionSource != "default" {
- log.Printf("[INFO] Handling NON-default execution source %s - NOT waiting or validating!", workflowExecution.ExecutionSource)
+ log.Printf("[INFO][%s] Handling NON-default execution source %s - NOT waiting or validating!", workflowExecution.ExecutionId, workflowExecution.ExecutionSource)
} else if workflowExecution.ExecutionSource == "default" {
- log.Printf("[INFO] Handling DEFAULT execution source %s - SKIPPING wait anyway due to exited issues!", workflowExecution.ExecutionSource)
+ log.Printf("[INFO][%s] Handling DEFAULT execution source %s - SKIPPING wait anyway due to exited issues!", workflowExecution.ExecutionId, workflowExecution.ExecutionSource)
}
- log.Printf("[DEBUG] Deployed container ID %s", cont.ID)
+ //log.Printf("[DEBUG] Deployed container ID %s", cont.ID)
//containerIds = append(containerIds, cont.ID)
return nil
@@ -620,11 +790,100 @@ func removeIndex(s []string, i int) []string {
return s[:len(s)-1]
}
+func getWorkerURLs() ([]string, error) {
+ workerUrls := []string{}
+
+ // Create a new Docker client
+ cli, err := dockerclient.NewEnvClient()
+ if err != nil {
+ log.Println("[ERROR] Failed to create Docker client:", err)
+ return workerUrls, err
+ }
+
+ // Specify the name of the service for which you want to list tasks
+ serviceName := "shuffle-workers"
+
+ // Get the list of tasks for the service
+ tasks, err := cli.TaskList(context.Background(), types.TaskListOptions{
+ Filters: filters.NewArgs(filters.Arg("service", serviceName)),
+ })
+
+ if err != nil {
+ log.Println("[ERROR] Failed to list tasks for service:", err)
+ return workerUrls, err
+ }
+
+ // Print task information
+ for _, task := range tasks {
+ url := fmt.Sprintf("http://%s.%d.%s:33333", serviceName, task.Slot, task.ID)
+ workerUrls = append(workerUrls, url)
+ }
+
+ return workerUrls, nil
+}
+
+func askOtherWorkersToDownloadImage(image string) {
+ if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
+ return
+ }
+
+ urls, err := getWorkerURLs()
+ if err != nil {
+ log.Printf("[ERROR] Error in listing worker urls: %s", err)
+ return
+ }
+
+ for _, url := range urls {
+ log.Printf("[DEBUG] Trying to speak to: %s", url)
+ imagesRequest := ImageRequest{
+ Image: image,
+ }
+
+ url = fmt.Sprintf("%s/api/v1/download", url)
+
+ imageJSON, err := json.Marshal(imagesRequest)
+
+ log.Printf("[INFO] Making a request to %s to download images", url)
+ req, err := http.NewRequest(
+ "POST",
+ url,
+ bytes.NewBuffer(imageJSON),
+ )
+
+ if err != nil {
+ log.Printf("[ERROR] Error in making request to %s : %s", url, err)
+ continue
+ }
+
+ httpClient := &http.Client{}
+ resp, err := httpClient.Do(req)
+ if err != nil {
+ log.Printf("[ERROR] Error in making request to %s : %s", url, err)
+ continue
+ }
+
+ defer resp.Body.Close()
+ respBody, err := ioutil.ReadAll(resp.Body)
+ if err != nil {
+ log.Printf("[ERROR] Error in reading response body : %s", err)
+ continue
+ }
+
+ log.Printf("[INFO] Response body when tried sending images for nodes to download: %s", respBody)
+ }
+}
+
func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
ctx := context.Background()
- //log.Printf("[DEBUG][%s] Pre DecideExecution", workflowExecution.ExecutionId)
workflowExecution, relevantActions := shuffle.DecideExecution(ctx, workflowExecution, environment)
+ if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "FAILURE" || workflowExecution.Status == "ABORTED" {
+ log.Printf("[DEBUG][%s] Shutting down because status is %s", workflowExecution.ExecutionId, workflowExecution.Status)
+ shutdown(workflowExecution, "", "Workflow run is already finished", true)
+ return
+ }
+
+
startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
dockercli, err := dockerclient.NewEnvClient()
@@ -633,7 +892,6 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
return
}
- // log.Printf("\n\n[DEBUG] Got %d relevant action(s) to run!\n\n", len(relevantActions))
for _, action := range relevantActions {
appname := action.AppName
appversion := action.AppVersion
@@ -645,6 +903,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if strings.Contains(image, " ") {
image = strings.ReplaceAll(image, " ", "-")
}
+ askOtherWorkersToDownloadImage(image)
// Added UUID to identifier just in case
//identifier := fmt.Sprintf("%s_%s_%s_%s_%s", appname, appversion, action.ID, workflowExecution.ExecutionId, uuid.NewV4())
@@ -688,7 +947,9 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
// marshal action and put it in there rofl
- log.Printf("[INFO][%s] Time to execute %s (%s) with app %s:%s, function %s, env %s with %d parameters.", workflowExecution.ExecutionId, action.ID, action.Label, action.AppName, action.AppVersion, action.Name, action.Environment, len(action.Parameters))
+ //log.Printf("[INFO][%s] Time to execute %s (%s) with app %s:%s, function %s, env %s with %d parameters.", workflowExecution.ExecutionId, action.ID, action.Label, action.AppName, action.AppVersion, action.Name, action.Environment, len(action.Parameters))
+
+ log.Printf("[DEBUG][%s] Action: Send, Label: '%s', Action: '%s', Run status: %s, Extra=", workflowExecution.ExecutionId, action.Label, action.AppName, workflowExecution.Status)
actionData, err := json.Marshal(action)
if err != nil {
@@ -697,7 +958,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
}
if action.AppID == "0ca8887e-b4af-4e3e-887c-87e9d3bc3d3e" {
- log.Printf("[DEBUG] Should run filter: %#v\n\n", action)
+ log.Printf("[DEBUG] Should run filter: %#v", action)
runFilter(workflowExecution, action)
continue
}
@@ -718,7 +979,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
fmt.Sprintf("CALLBACK_URL=%s", baseUrl),
fmt.Sprintf("BASE_URL=%s", appCallbackUrl),
fmt.Sprintf("TZ=%s", timezone),
- fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", os.Getenv("SHUFFLE_LOGS_DISABLED")),
+ fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", logsDisabled),
}
if len(actionData) >= 100000 {
@@ -762,6 +1023,21 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
env = append(env, fmt.Sprintf("NO_PROXY=%s", os.Getenv("NO_PROXY")))
}
+ overrideHttpProxy := os.Getenv("SHUFFLE_INTERNAL_HTTP_PROXY")
+ overrideHttpsProxy := os.Getenv("SHUFFLE_INTERNAL_HTTPS_PROXY")
+ if overrideHttpProxy != "" {
+ env = append(env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTP_PROXY=%s", overrideHttpProxy))
+ }
+
+ if overrideHttpsProxy != "" {
+ env = append(env, fmt.Sprintf("SHUFFLE_INTERNAL_HTTPS_PROXY=%s", overrideHttpsProxy))
+ }
+
+ if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 {
+ env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
+ }
+
+
// Fixes issue:
// standard_go init_linux.go:185: exec user process caused "argument list too long"
// https://devblogs.microsoft.com/oldnewthing/20100203-00/?p=15083
@@ -785,10 +1061,12 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
// 3. Add remote repo location
images := []string{
image,
- fmt.Sprintf("%s:%s_%s", baseimagename, parsedAppname, action.AppVersion),
fmt.Sprintf("%s/%s:%s_%s", registryName, baseimagename, parsedAppname, action.AppVersion),
+ fmt.Sprintf("%s:%s_%s", baseimagename, parsedAppname, action.AppVersion),
}
+
+
// If cleanup is set, it should run for efficiency
pullOptions := types.ImagePullOptions{}
if cleanupEnv == "true" {
@@ -859,7 +1137,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
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("[ERROR] Docker build:%sERROR 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
@@ -969,7 +1247,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
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("[ERROR] Docker build:%sERROR 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
@@ -1018,7 +1296,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
// FIXME - clean up stopped (remove) containers with this execution id
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 (2)\n\n", workflowExecution.ExecutionId, err)
+ log.Printf("[ERROR] Failed to update exec variables for execution %s: %s (2)", workflowExecution.ExecutionId, err)
}
if len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions)+extra {
@@ -1036,13 +1314,22 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
if shutdownCheck {
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)
+ validated := shuffle.ValidateFinished(ctx, -1, workflowExecution)
+ if validated {
+ shutdownData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed marshalling shutdowndata during set: %s", err)
+ }
+
+ sendResult(workflowExecution, shutdownData)
+ }
+
log.Printf("[DEBUG][%s] Shutting down (17)", workflowExecution.ExecutionId)
- if os.Getenv("IS_KUBERNETES") == "true" {
+ if isKubernetes == "true" {
// log.Printf("workflow execution: %#v", workflowExecution)
clientset, err := getKubernetesClient()
if err != nil {
- fmt.Println("[ERROR]Error getting kubernetes client:", err)
+ log.Println("[ERROR] Error getting kubernetes client (1):", err)
os.Exit(1)
}
cleanupExecution(clientset, workflowExecution, "shuffle")
@@ -1094,7 +1381,7 @@ func executionInit(workflowExecution shuffle.WorkflowExecution) error {
}
for _, trigger := range workflowExecution.Workflow.Triggers {
- //log.Printf("Appname trigger (0): %s", trigger.AppName)
+ //log.Printf("Appname trigger (0): %s (%s)", trigger.AppName, trigger.ID)
if trigger.AppName == "User Input" || trigger.AppName == "Shuffle Workflow" {
if trigger.ID == branch.SourceID {
sourceFound = true
@@ -1107,22 +1394,16 @@ func executionInit(workflowExecution shuffle.WorkflowExecution) error {
if sourceFound {
parents[branch.DestinationID] = append(parents[branch.DestinationID], branch.SourceID)
} else {
- log.Printf("[DEBUG] ID %s was not found in actions! Skipping parent. (TRIGGER?)", branch.SourceID)
+ log.Printf("[DEBUG] Parent ID %s was not found in actions! Skipping parent. (TRIGGER?)", branch.SourceID)
}
if destinationFound {
children[branch.SourceID] = append(children[branch.SourceID], branch.DestinationID)
} else {
- log.Printf("[DEBUG] ID %s was not found in actions! Skipping child. (TRIGGER?)", branch.SourceID)
+ log.Printf("[DEBUG] Child ID %s was not found in actions! Skipping child. (TRIGGER?)", branch.SourceID)
}
}
- /*
- log.Printf("\n\n\n[INFO] CHILDREN FOUND: %#v", children)
- log.Printf("[INFO] PARENTS FOUND: %#v", parents)
- log.Printf("[INFO] NEXT ACTIONS: %#v\n\n", nextActions)
- */
-
log.Printf("[INFO][%s] shuffle.Actions: %d + Special shuffle.Triggers: %d", workflowExecution.ExecutionId, len(workflowExecution.Workflow.Actions), extra)
onpremApps := []string{}
toExecuteOnprem := []string{}
@@ -1190,16 +1471,185 @@ func executionInit(workflowExecution shuffle.WorkflowExecution) error {
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)
+ log.Printf("[ERROR] Failed to update exec variables for execution %s: %s", workflowExecution.ExecutionId, err)
}
return nil
}
+func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.WorkflowExecution, streamResultUrl, subflowId string) error {
+ extra := 0
+ for _, trigger := range workflowExecution.Workflow.Triggers {
+ if trigger.AppName == "User Input" || trigger.AppName == "Shuffle Workflow" {
+ extra += 1
+ }
+ }
+
+ req, err := http.NewRequest(
+ "POST",
+ streamResultUrl,
+ bytes.NewBuffer([]byte(data)),
+ )
+
+ newresp, err := topClient.Do(req)
+ if err != nil {
+ log.Printf("[ERROR] Failed making request (1): %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ defer newresp.Body.Close()
+ body, err := ioutil.ReadAll(newresp.Body)
+ if err != nil {
+ log.Printf("[ERROR] Failed reading body (1): %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ if newresp.StatusCode != 200 {
+ 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)
+ }
+
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return errors.New(fmt.Sprintf("Bad statuscode: %d", newresp.StatusCode))
+ }
+
+ err = json.Unmarshal(body, &workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
+ log.Printf("[INFO][%s] Workflow execution is finished. Exiting worker.", workflowExecution.ExecutionId)
+ log.Printf("[DEBUG] Shutting down (20)")
+ if isKubernetes == "true" {
+ // log.Printf("workflow execution: %#v", workflowExecution)
+ clientset, err := getKubernetesClient()
+ if err != nil {
+ log.Println("[ERROR] Error getting kubernetes client (2):", err)
+ os.Exit(1)
+ }
+
+ cleanupExecution(clientset, workflowExecution, "shuffle")
+ } else {
+ shutdown(workflowExecution, "", "", true)
+ }
+ }
+
+ for _, result := range workflowExecution.Results {
+ if result.Action.ID != subflowId {
+ continue
+ }
+
+ log.Printf("[DEBUG][%s] Found subflow to handle: %s (%s)", workflowExecution.ExecutionId, result.Action.Label, result.Status)
+ if result.Status == "SUCCESS" || result.Status == "FINISHED" || result.Status == "FAILURE" || result.Status == "ABORTED" {
+ // Check for results
+
+ setWorkflowExecution(ctx, workflowExecution, false)
+ return nil
+ }
+ }
+
+ log.Printf("[INFO][%s] Status: %s, Results: %d, actions: %d", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
+ return errors.New("Subflow status not found yet")
+}
+
+func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffle.WorkflowExecution, streamResultUrl string, extra int) error {
+ if extra == -1 {
+ extra = 0
+ for _, trigger := range workflowExecution.Workflow.Triggers {
+ if trigger.AppName == "User Input" || trigger.AppName == "Shuffle Workflow" {
+ extra += 1
+ }
+ }
+ }
+
+ req, err := http.NewRequest(
+ "POST",
+ streamResultUrl,
+ bytes.NewBuffer([]byte(data)),
+ )
+
+ newresp, err := topClient.Do(req)
+ if err != nil {
+ log.Printf("[ERROR] Failed making request (1): %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ defer newresp.Body.Close()
+ body, err := ioutil.ReadAll(newresp.Body)
+ if err != nil {
+ log.Printf("[ERROR] Failed reading body (1): %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ if newresp.StatusCode != 200 {
+ 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)
+ }
+
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return errors.New(fmt.Sprintf("Bad statuscode: %d", newresp.StatusCode))
+ }
+
+ err = json.Unmarshal(body, &workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return err
+ }
+
+ if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
+ log.Printf("[INFO][%s] Workflow execution is finished. Exiting worker.", workflowExecution.ExecutionId)
+ log.Printf("[DEBUG] Shutting down (20)")
+ if isKubernetes == "true" {
+ // log.Printf("workflow execution: %#v", workflowExecution)
+ clientset, err := getKubernetesClient()
+ if err != nil {
+ log.Println("[ERROR] Error getting kubernetes client (2):", err)
+ os.Exit(1)
+ }
+ cleanupExecution(clientset, workflowExecution, "shuffle")
+ } else {
+ shutdown(workflowExecution, "", "", true)
+ }
+ }
+
+ log.Printf("[INFO][%s] Status: %s, Results: %d, actions: %d", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
+ if workflowExecution.Status != "EXECUTING" {
+ log.Printf("[WARNING][%s] Exiting as worker execution has status %s!", workflowExecution.ExecutionId, workflowExecution.Status)
+ log.Printf("[DEBUG] Shutting down (21)")
+ if isKubernetes == "true" {
+ // log.Printf("workflow execution: %#v", workflowExecution)
+ clientset, err := getKubernetesClient()
+ if err != nil {
+ log.Println("[ERROR] Error getting kubernetes client (3):", err)
+ os.Exit(1)
+ }
+ cleanupExecution(clientset, workflowExecution, "shuffle")
+ } else {
+ shutdown(workflowExecution, "", "", true)
+ }
+ }
+
+ setWorkflowExecution(ctx, workflowExecution, false)
+ 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
@@ -1220,84 +1670,10 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl)
for {
- //fullUrl := fmt.Sprintf("%s/api/v1/workflows/%s/executions/%s/abort", baseUrl, workflowExecution.Workflow.ID, workflowExecution.ExecutionId)
- //log.Printf("[INFO] URL: %s", fullUrl)
- req, err := http.NewRequest(
- "POST",
- streamResultUrl,
- bytes.NewBuffer([]byte(data)),
- )
-
- newresp, err := topClient.Do(req)
+ err = handleDefaultExecutionWrapper(ctx, workflowExecution, streamResultUrl, extra)
if err != nil {
- log.Printf("[ERROR] Failed making request (1): %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
+ log.Printf("[ERROR] Failed handling default execution: %s", err)
}
-
- defer newresp.Body.Close()
- body, err := ioutil.ReadAll(newresp.Body)
- if err != nil {
- log.Printf("[ERROR] Failed reading body (1): %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- if newresp.StatusCode != 200 {
- 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)
- }
-
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- err = json.Unmarshal(body, &workflowExecution)
- if err != nil {
- log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
- log.Printf("[INFO][%s] Workflow execution is finished. Exiting worker.", workflowExecution.ExecutionId)
- log.Printf("[DEBUG] Shutting down (20)")
- //handle workerssssssssss
- if os.Getenv("IS_KUBERNETES") == "true" {
- // log.Printf("workflow execution: %#v", workflowExecution)
- clientset, err := getKubernetesClient()
- if err != nil {
- fmt.Println("[ERROR]Error getting kubernetes client:", err)
- os.Exit(1)
- }
- cleanupExecution(clientset, workflowExecution, "shuffle")
- } else {
- shutdown(workflowExecution, "", "", true)
- }
- }
-
- log.Printf("[INFO][%s] Status: %s, Results: %d, actions: %d", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra)
- if workflowExecution.Status != "EXECUTING" {
- log.Printf("[WARNING][%s] Exiting as worker execution has status %s!", workflowExecution.ExecutionId, workflowExecution.Status)
- log.Printf("[DEBUG] Shutting down (21)")
- if os.Getenv("IS_KUBERNETES") == "true" {
- // log.Printf("workflow execution: %#v", workflowExecution)
- clientset, err := getKubernetesClient()
- if err != nil {
- fmt.Println("[ERROR]Error getting kubernetes client:", err)
- os.Exit(1)
- }
- cleanupExecution(clientset, workflowExecution, "shuffle")
- } else {
- shutdown(workflowExecution, "", "", true)
- }
- }
-
- setWorkflowExecution(ctx, workflowExecution, false)
- //handleExecutionResult(workflowExecution)
}
return nil
@@ -1377,7 +1753,7 @@ func runSkipAction(client *http.Client, action shuffle.Action, workflowId, workf
return err
}
- newresp, err := client.Do(req)
+ newresp, err := topClient.Do(req)
if err != nil {
log.Printf("[WARNING] Error running skip request (0): %s", err)
return err
@@ -1394,167 +1770,6 @@ func runSkipAction(client *http.Client, action shuffle.Action, workflowId, workf
return nil
}
-// Sends request back to backend to handle the node
-func runUserInput(client *http.Client, action shuffle.Action, workflowId string, workflowExecution shuffle.WorkflowExecution, authorization string, configuration string, dockercli *dockerclient.Client) error {
- timeNow := time.Now().Unix()
- result := shuffle.ActionResult{
- Action: action,
- ExecutionId: workflowExecution.ExecutionId,
- Authorization: authorization,
- Result: configuration,
- StartedAt: timeNow,
- CompletedAt: 0,
- Status: "WAITING",
- }
-
- // Checking for userinput to deploy subflow for it
- subflow := false
- subflowId := ""
- argument := ""
- continueUrl := "testing continue"
- cancelUrl := "testing cancel"
- for _, item := range action.Parameters {
- if item.Name == "subflow" {
- subflow = true
- subflowId = item.Value
- } else if item.Name == "alertinfo" {
- argument = item.Value
- }
- }
-
- if subflow {
- log.Printf("[DEBUG] Should run action with subflow app with argument %#v", argument)
- newAction := shuffle.Action{
- AppName: "shuffle-subflow",
- Name: "run_subflow",
- AppVersion: "1.0.0",
- Label: "User Input Subflow Execution",
- }
-
- identifier := fmt.Sprintf("%s_%s_%s_%s", newAction.AppName, newAction.AppVersion, action.ID, workflowExecution.ExecutionId)
- if strings.Contains(identifier, " ") {
- identifier = strings.ReplaceAll(identifier, " ", "-")
- }
-
- inputValue := UserInputSubflow{
- Argument: argument,
- ContinueUrl: continueUrl,
- CancelUrl: cancelUrl,
- }
-
- parsedArgument, err := json.Marshal(inputValue)
- if err != nil {
- log.Printf("[ERROR] Failed to parse arguments: %s", err)
- parsedArgument = []byte(argument)
- }
-
- newAction.Parameters = []shuffle.WorkflowAppActionParameter{
- shuffle.WorkflowAppActionParameter{
- Name: "user_apikey",
- Value: workflowExecution.Authorization,
- },
- shuffle.WorkflowAppActionParameter{
- Name: "workflow",
- Value: subflowId,
- },
- shuffle.WorkflowAppActionParameter{
- Name: "argument",
- Value: string(parsedArgument),
- },
- }
-
- newAction.Parameters = append(newAction.Parameters, shuffle.WorkflowAppActionParameter{
- Name: "source_workflow",
- Value: workflowExecution.Workflow.ID,
- })
-
- newAction.Parameters = append(newAction.Parameters, shuffle.WorkflowAppActionParameter{
- Name: "source_execution",
- Value: workflowExecution.ExecutionId,
- })
-
- newAction.Parameters = append(newAction.Parameters, shuffle.WorkflowAppActionParameter{
- Name: "source_node",
- Value: action.ID,
- })
-
- newAction.Parameters = append(newAction.Parameters, shuffle.WorkflowAppActionParameter{
- Name: "source_auth",
- Value: workflowExecution.Authorization,
- })
-
- newAction.Parameters = append(newAction.Parameters, shuffle.WorkflowAppActionParameter{
- Name: "startnode",
- Value: "",
- })
-
- // If cleanup is set, it should run for efficiency
- //appName := strings.Replace(identifier, fmt.Sprintf("_%s", action.ID), "", -1)
- //appName = strings.Replace(appName, fmt.Sprintf("_%s", workflowExecution.ExecutionId), "", -1)
- actionData, err := json.Marshal(newAction)
- if err != nil {
- return err
- }
-
- 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", baseUrl),
- fmt.Sprintf("BASE_URL=%s", appCallbackUrl),
- fmt.Sprintf("TZ=%s", timezone),
- fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", os.Getenv("SHUFFLE_LOGS_DISABLED")),
- }
-
- if strings.ToLower(os.Getenv("SHUFFLE_PASS_APP_PROXY")) == "true" {
- //log.Printf("APPENDING PROXY TO THE APP!")
- env = append(env, fmt.Sprintf("HTTP_PROXY=%s", os.Getenv("HTTP_PROXY")))
- env = append(env, fmt.Sprintf("HTTPS_PROXY=%s", os.Getenv("HTTPS_PROXY")))
- env = append(env, fmt.Sprintf("NO_PROXY=%s", os.Getenv("NO_PROXY")))
- }
-
- err = deployApp(dockercli, "frikky/shuffle:shuffle-subflow_1.0.0", identifier, env, workflowExecution, newAction)
- if err != nil {
- log.Printf("[ERROR] Failed to deploy subflow for user input trigger %s: %s", action.ID, err)
- }
- } else {
- log.Printf("[DEBUG] Running user input WITHOUT subflow")
- }
-
- resultData, err := json.Marshal(result)
- if err != nil {
- return err
- }
-
- streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl)
- req, err := http.NewRequest(
- "POST",
- streamUrl,
- bytes.NewBuffer([]byte(resultData)),
- )
-
- if err != nil {
- log.Printf("[WARNING] Error building test request (2): %s", err)
- return err
- }
-
- newresp, err := client.Do(req)
- if err != nil {
- log.Printf("[WARNING] Error running test request (2): %s", err)
- return err
- }
-
- defer newresp.Body.Close()
- body, err := ioutil.ReadAll(newresp.Body)
- if err != nil {
- log.Printf("Failed reading body when waiting: %s", err)
- return err
- }
-
- log.Printf("[INFO] User Input Body: %s", string(body))
- return nil
-}
-
func runTestExecution(client *http.Client, workflowId, apikey string) (string, string) {
executeUrl := fmt.Sprintf("%s/api/v1/workflows/%s/execute", baseUrl, workflowId)
req, err := http.NewRequest(
@@ -1569,7 +1784,7 @@ func runTestExecution(client *http.Client, workflowId, apikey string) (string, s
}
req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", apikey))
- newresp, err := client.Do(req)
+ newresp, err := topClient.Do(req)
if err != nil {
log.Printf("[WARNING] Error running test request (3): %s", err)
return "", ""
@@ -1593,12 +1808,53 @@ func runTestExecution(client *http.Client, workflowId, apikey string) (string, s
return workflowExecution.Authorization, workflowExecution.ExecutionId
}
+func isRunningInCluster() bool {
+ _, existsHost := os.LookupEnv("KUBERNETES_SERVICE_HOST")
+ _, existsPort := os.LookupEnv("KUBERNETES_SERVICE_PORT")
+ return existsHost && existsPort
+}
+
+func buildEnvVars(envMap map[string]string) []corev1.EnvVar {
+ var envVars []corev1.EnvVar
+ for key, value := range envMap {
+ envVars = append(envVars, corev1.EnvVar{Name: key, Value: value})
+ }
+ return envVars
+}
+
+func getKubernetesClient() (*kubernetes.Clientset, error) {
+ if isRunningInCluster() {
+ config, err := rest.InClusterConfig()
+ if err != nil {
+ return nil, err
+ }
+ clientset, err := kubernetes.NewForConfig(config)
+ if err != nil {
+ return nil, err
+ }
+ return clientset, nil
+ } else {
+ home := homedir.HomeDir()
+ kubeconfigPath := filepath.Join(home, ".kube", "config")
+ config, err := clientcmd.BuildConfigFromFlags("", kubeconfigPath)
+ if err != nil {
+ return nil, err
+ }
+ clientset, err := kubernetes.NewForConfig(config)
+ if err != nil {
+ return nil, err
+ }
+ return clientset, nil
+ }
+}
+
func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
if request.Body == nil {
resp.WriteHeader(http.StatusBadRequest)
return
}
+ defer request.Body.Close()
body, err := ioutil.ReadAll(request.Body)
if err != nil {
log.Printf("[WARNING] (3) Failed reading body for workflowqueue")
@@ -1607,8 +1863,6 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
return
}
- defer request.Body.Close()
-
var actionResult shuffle.ActionResult
err = json.Unmarshal(body, &actionResult)
if err != nil {
@@ -1619,7 +1873,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
}
if len(actionResult.ExecutionId) == 0 {
- log.Printf("[WARNING] No workflow execution id in action result. Data: %s", string(body))
+ log.Printf("[ERROR] No workflow execution id in action result. Data: %s", string(body))
resp.WriteHeader(400)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "No workflow execution id in action result"}`)))
return
@@ -1635,44 +1889,47 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
workflowExecution, err := shuffle.GetWorkflowExecution(ctx, actionResult.ExecutionId)
if err != nil {
log.Printf("[ERROR][%s] Failed getting execution (workflowqueue) %s: %s", actionResult.ExecutionId, actionResult.ExecutionId, err)
- resp.WriteHeader(401)
+ resp.WriteHeader(500)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting execution ID %s because it doesn't exist locally."}`, actionResult.ExecutionId)))
return
}
if workflowExecution.Authorization != actionResult.Authorization {
- log.Printf("[INFO] Bad authorization key when updating node (workflowQueue) %s. Want: %s, Have: %s", actionResult.ExecutionId, workflowExecution.Authorization, actionResult.Authorization)
- resp.WriteHeader(401)
+ log.Printf("[ERROR][%s] Bad authorization key when updating node (workflowQueue). Want: %s, Have: %s", actionResult.ExecutionId, workflowExecution.Authorization, actionResult.Authorization)
+ resp.WriteHeader(403)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key"}`)))
return
}
if workflowExecution.Status == "FINISHED" {
- log.Printf("[DEBUG] Workflowexecution is already FINISHED. No further action can be taken")
- resp.WriteHeader(401)
+ log.Printf("[DEBUG][%s] Workflowexecution is already FINISHED. No further action can be taken", workflowExecution.ExecutionId)
+ resp.WriteHeader(200)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Workflowexecution is already finished because it has status %s. Lastnode: %s"}`, workflowExecution.Status, workflowExecution.LastNode)))
return
}
if workflowExecution.Status == "ABORTED" || workflowExecution.Status == "FAILURE" {
+ log.Printf("[WARNING][%s] Workflowexecution already has status %s. No further action can be taken", workflowExecution.ExecutionId, workflowExecution.Status)
+ resp.WriteHeader(200)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Workflowexecution is aborted because of %s with result %s and status %s"}`, workflowExecution.LastNode, workflowExecution.Result, workflowExecution.Status)))
+ return
+ }
- if workflowExecution.Workflow.Configuration.ExitOnError {
- log.Printf("[WARNING] Workflowexecution already has status %s. No further action can be taken", workflowExecution.Status)
- resp.WriteHeader(401)
- resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Workflowexecution is aborted because of %s with result %s and status %s"}`, workflowExecution.LastNode, workflowExecution.Result, workflowExecution.Status)))
- return
- } else {
- log.Printf("Continuing even though it's aborted.")
+ retries := 0
+ retry, retriesok := request.URL.Query()["retries"]
+ if retriesok && len(retry) > 0 {
+ val, err := strconv.Atoi(retry[0])
+ if err == nil {
+ retries = val
}
}
- log.Printf("[INFO][%s] Got result '%s' from '%s' with app '%s':'%s'", actionResult.ExecutionId, actionResult.Status, actionResult.Action.Label, actionResult.Action.AppName, actionResult.Action.AppVersion)
+ log.Printf("[DEBUG][%s] Action: Received, Label: '%s', Action: '%s', Status: %s, Run status: %s, Extra=Retry:%d", workflowExecution.ExecutionId, actionResult.Action.Label, actionResult.Action.AppName, actionResult.Status, workflowExecution.Status, retries)
//results = append(results, actionResult)
//log.Printf("[INFO][%s] Time to execute %s (%s) with app %s:%s, function %s, env %s with %d parameters.", workflowExecution.ExecutionId, action.ID, action.Label, action.AppName, action.AppVersion, action.Name, action.Environment, len(action.Parameters))
//log.Printf("[DEBUG][%s] In workflowQueue with transaction", workflowExecution.ExecutionId)
runWorkflowExecutionTransaction(ctx, 0, workflowExecution.ExecutionId, actionResult, resp)
-
}
// Will make sure transactions are always ran for an execution. This is recursive if it fails. Allowed to fail up to 5 times
@@ -1690,12 +1947,38 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
setExecution := true
workflowExecution, dbSave, err := shuffle.ParsedExecutionResult(ctx, *workflowExecution, actionResult, true, 0)
- if err != nil {
+ if err == nil {
+ if workflowExecution.Status != "EXECUTING" && workflowExecution.Status != "WAITING" {
+ log.Printf("[WARNING][%s] Execution is not executing, but %s. Stopping Transaction update.", workflowExecution.ExecutionId, workflowExecution.Status)
+ if resp != nil {
+ resp.WriteHeader(200)
+ resp.Write([]byte(fmt.Sprintf(`{"success": true, "reason": "Execution is not executing, but %s"}`, workflowExecution.Status)))
+ }
+
+
+ log.Printf("[DEBUG][%s] Shutting down (35)", workflowExecution.ExecutionId)
+
+ // Force sending result
+ shutdownData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR][%s] Failed marshalling execution (35): %s", workflowExecution.ExecutionId, err)
+ }
+
+ sendResult(*workflowExecution, shutdownData)
+ shutdown(*workflowExecution, "", "", false)
+ return
+ }
+ } else {
+ if strings.Contains(strings.ToLower(fmt.Sprintf("%s", err)), "already been ran") || strings.Contains(strings.ToLower(fmt.Sprintf("%s", err)), "already finished") {
+ log.Printf("[ERROR][%s] Skipping rerun of action result as it's already been ran: %s", workflowExecution.ExecutionId)
+ return
+ }
+
log.Printf("[DEBUG] Rerunning transaction? %s", err)
if strings.Contains(fmt.Sprintf("%s", err), "Rerun this transaction") {
workflowExecution, err := shuffle.GetWorkflowExecution(ctx, workflowExecutionId)
if err != nil {
- log.Printf("[ERROR] Failed getting execution cache (2): %s", err)
+ log.Printf("[ERROR][%s] Failed getting execution cache (2): %s", workflowExecution.ExecutionId, err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting execution (2)"}`)))
return
@@ -1706,15 +1989,15 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
workflowExecution, dbSave, err = shuffle.ParsedExecutionResult(ctx, *workflowExecution, actionResult, false, 0)
if err != nil {
- log.Printf("[ERROR] Failed execution of parsedexecution (2): %s", err)
+ log.Printf("[ERROR][%s] Failed execution of parsedexecution (2): %s", workflowExecution.ExecutionId, err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting execution (2)"}`)))
return
} else {
- log.Printf("[DEBUG] Successfully got ParsedExecution with %d results!", len(workflowExecution.Results))
+ log.Printf("[DEBUG][%s] Successfully got ParsedExecution with %d results!", workflowExecution.ExecutionId, len(workflowExecution.Results))
}
} else {
- log.Printf("[ERROR] Failed execution of parsedexecution: %s", err)
+ log.Printf("[ERROR][%s] Failed execution of parsedexecution: %s", workflowExecution.ExecutionId, err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting execution"}`)))
return
@@ -1739,43 +2022,28 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
cacheData := []byte(cache.([]uint8))
err = json.Unmarshal(cacheData, &workflowExecution)
if err != nil {
- log.Printf("[ERROR] Failed unmarshalling workflowexecution: %s", err)
+ log.Printf("[ERROR][%s] Failed unmarshalling workflowexecution: %s", workflowExecution.ExecutionId, err)
}
if len(parsedValue.Results) > 0 && len(parsedValue.Results) != resultLength {
setExecution = false
if attempts > 5 {
- //log.Printf("\n\nSkipping execution input - %d vs %d. Attempts: (%d)\n\n", len(parsedValue.Results), resultLength, attempts)
}
attempts += 1
+ log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(parsedValue.Results), resultLength)
+ /*
if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) {
+ log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions))
runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp)
return
}
+ */
}
}
- /*
- if value, found := requestCache.Get(cacheKey); found {
- parsedValue := value.(*shuffle.WorkflowExecution)
- if len(parsedValue.Results) > 0 && len(parsedValue.Results) != resultLength {
- setExecution = false
- if attempts > 5 {
- //log.Printf("\n\nSkipping execution input - %d vs %d. Attempts: (%d)\n\n", len(parsedValue.Results), resultLength, attempts)
- }
-
- attempts += 1
- if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) {
- runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp)
- return
- }
- }
- }
- */
-
if setExecution || workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || workflowExecution.Status == "FAILURE" {
- log.Printf("[DEBUG][%s] Running setexec with status %s and %d results", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results))
+ log.Printf("[DEBUG][%s] Running setexec with status %s and %d result(s)", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results))
err = setWorkflowExecution(ctx, *workflowExecution, dbSave)
if err != nil {
resp.WriteHeader(401)
@@ -1789,7 +2057,6 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
// Just in case. Should MAYBE validate finishing another time as well.
// This fixes issues with e.g. shuffle.Action -> shuffle.Trigger -> shuffle.Action.
handleExecutionResult(*workflowExecution)
- //validateFinished(workflowExecution)
}
//if newExecutions && len(nextActions) > 0 {
@@ -1802,8 +2069,7 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
}
func sendSelfRequest(actionResult shuffle.ActionResult) {
- log.Printf("[INFO][%s] Not sending backend info since source is default (not swarm)", actionResult.ExecutionId)
- return
+
data, err := json.Marshal(actionResult)
if err != nil {
@@ -1837,7 +2103,7 @@ func sendSelfRequest(actionResult shuffle.ActionResult) {
newresp, err := topClient.Do(req)
if err != nil {
- log.Printf("[ERROR][%s] Error running self request (2): %s", actionResult.ExecutionId, err)
+ log.Printf("[ERROR][%s] Error running finishing request (2): %s", actionResult.ExecutionId, err)
return
}
@@ -1846,9 +2112,9 @@ func sendSelfRequest(actionResult shuffle.ActionResult) {
body, err := ioutil.ReadAll(newresp.Body)
//log.Printf("[INFO] BACKEND STATUS: %d", newresp.StatusCode)
if err != nil {
- log.Printf("[ERROR][%s] Failed reading self request body: %s", actionResult.ExecutionId, err)
+ log.Printf("[ERROR][%s] Failed reading body: %s", actionResult.ExecutionId, err)
} else {
- log.Printf("[DEBUG][%s] NEWRESP (from self - 1): %s", actionResult.ExecutionId, string(body))
+ log.Printf("[DEBUG][%s] NEWRESP (from backend): %s", actionResult.ExecutionId, string(body))
}
}
}
@@ -1857,8 +2123,27 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
if workflowExecution.ExecutionSource == "default" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
//log.Printf("[INFO][%s] Not sending backend info since source is default (not swarm)", workflowExecution.ExecutionId)
//return
+ } else {
}
+ // Basically to reduce backend strain
+ /*
+ if shuffle.ArrayContains(finishedExecutions, workflowExecution.ExecutionId) {
+ log.Printf("[INFO][%s] NOT sending backend info since it's already been sent before.", workflowExecution.ExecutionId)
+ return
+ }
+ */
+
+ // Take it down again
+ /*
+ if len(finishedExecutions) > 100 {
+ log.Printf("[DEBUG][%s] Removing old execution from finishedExecutions: %s", workflowExecution.ExecutionId, finishedExecutions[0])
+ finishedExecutions = finishedExecutions[99:]
+ }
+
+ finishedExecutions = append(finishedExecutions, workflowExecution.ExecutionId)
+ */
+
streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl)
req, err := http.NewRequest(
"POST",
@@ -1907,10 +2192,9 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) bool {
workflowExecution = shuffle.Fixexecution(ctx, workflowExecution)
_, extra, _, _, _, _, _, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
- log.Printf("[INFO][%s] VALIDATION. Status: %s, shuffle.Actions: %d, Extra: %d, Results: %d. Parent: %#v\n", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Workflow.Actions), extra, len(workflowExecution.Results), workflowExecution.ExecutionParent)
+ log.Printf("[INFO][%s] VALIDATION. Status: %s, shuffle.Actions: %d, Extra: %d, Results: %d. Parent: %#v", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Workflow.Actions), extra, len(workflowExecution.Results), workflowExecution.ExecutionParent)
- //if len(workflowExecution.Results) == len(workflowExecution.Workflow.Actions)+extra {
- if (len(environments) == 1 && requestsSent == 0 && len(workflowExecution.Results) >= 1 && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm") || (len(workflowExecution.Results) >= len(workflowExecution.Workflow.Actions)+extra && len(workflowExecution.Workflow.Actions) > 0) {
+ if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || (len(environments) == 1 && requestsSent == 0 && len(workflowExecution.Results) >= 1 && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm") || (len(workflowExecution.Results) >= len(workflowExecution.Workflow.Actions)+extra && len(workflowExecution.Workflow.Actions) > 0) {
if workflowExecution.Status == "FINISHED" {
for _, result := range workflowExecution.Results {
@@ -1921,7 +2205,6 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) bool {
}
}
- requestsSent += 1
log.Printf("[DEBUG][%s] Should send full result to %s", workflowExecution.ExecutionId, baseUrl)
@@ -1933,9 +2216,11 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) bool {
}
cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId)
- err = shuffle.SetCache(ctx, cacheKey, shutdownData, 30)
- if err != nil {
- log.Printf("[ERROR][%s] Failed adding to cache during validateFinished", workflowExecution)
+ if len(workflowExecution.Authorization) > 0 {
+ err = shuffle.SetCache(ctx, cacheKey, shutdownData, 31)
+ if err != nil {
+ log.Printf("[ERROR][%s] Failed adding to cache during ValidateFinished", workflowExecution)
+ }
}
shuffle.RunCacheCleanup(ctx, workflowExecution)
@@ -1947,6 +2232,7 @@ func validateFinished(workflowExecution shuffle.WorkflowExecution) bool {
}
func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
+ defer request.Body.Close()
body, err := ioutil.ReadAll(request.Body)
if err != nil {
log.Printf("[WARNING] Failed reading body for stream result queue")
@@ -1955,9 +2241,6 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
return
}
- defer request.Body.Close()
- //log.Printf("[DEBUG] In get stream results with body length %d: %s", len(body), string(body))
-
var actionResult shuffle.ActionResult
err = json.Unmarshal(body, &actionResult)
if err != nil {
@@ -1985,7 +2268,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
// Authorization is done here
if workflowExecution.Authorization != actionResult.Authorization {
- log.Printf("Bad authorization key when getting stream results %s.", actionResult.ExecutionId)
+ log.Printf("[ERROR] Bad authorization key when getting stream results from cache %s.", actionResult.ExecutionId)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key or execution_id might not exist."}`)))
return
@@ -2003,52 +2286,10 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
}
-func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.WorkflowExecution, dbSave bool) error {
- if len(workflowExecution.ExecutionId) == 0 {
- log.Printf("[DEBUG] Workflowexecution executionId can't be empty.")
- return errors.New("ExecutionId can't be empty.")
- }
-
- //log.Printf("[DEBUG][%s] Setting with %d results (pre)", workflowExecution.ExecutionId, len(workflowExecution.Results))
- workflowExecution = shuffle.Fixexecution(ctx, workflowExecution)
- //log.Printf("[DEBUG][%s] Setting with %d results (post)", workflowExecution.ExecutionId, len(workflowExecution.Results))
-
- cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId)
-
- execData, err := json.Marshal(workflowExecution)
- if err != nil {
- log.Printf("[ERROR] Failed marshalling execution during set: %s", err)
- return err
- }
-
- err = shuffle.SetCache(ctx, cacheKey, execData, 30)
- if err != nil {
- log.Printf("[ERROR][%s] Failed adding to cache during setexecution", workflowExecution)
- return err
- }
- //requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration)
-
- handleExecutionResult(workflowExecution)
- validateFinished(workflowExecution)
-
- // FIXME: Should this shutdown OR send the result?
- // The worker may not be running the backend hmm
- if dbSave {
- if workflowExecution.ExecutionSource == "default" {
- log.Printf("[DEBUG][%s] Shutting down (25)", workflowExecution.ExecutionId)
- shutdown(workflowExecution, "", "", true)
- //return
- } else {
- log.Printf("[DEBUG] NOT shutting down with dbSave (%s)", workflowExecution.ExecutionSource)
- }
- }
-
- return nil
-}
-
// GetLocalIP returns the non loopback local IP of the host
func getLocalIP() string {
+
addrs, err := net.InterfaceAddrs()
if err != nil {
return ""
@@ -2082,7 +2323,6 @@ func getAvailablePort() (net.Listener, error) {
func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
hostname = getLocalIP()
-
os.Setenv("WORKER_HOSTNAME", hostname)
// FIXME: This MAY not work because of speed between first
@@ -2094,9 +2334,11 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
}
log.Printf("[DEBUG] OLD HOSTNAME: %s", appCallbackUrl)
+
+
port := listener.Addr().(*net.TCPAddr).Port
- log.Printf("\n\n[DEBUG] Starting webserver (2) on port %d with hostname: %s\n\n", port, hostname)
+ log.Printf("[DEBUG] Starting webserver (2) on port %d with hostname: %s", port, hostname)
appCallbackUrl = fmt.Sprintf("http://%s:%d", hostname, port)
log.Printf("[INFO] NEW WORKER HOSTNAME: %s", appCallbackUrl)
@@ -2130,7 +2372,7 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error {
//return
}
- newresp, err := client.Do(req)
+ newresp, err := topClient.Do(req)
if err != nil {
log.Printf("[ERROR] Failed download request for %s: %s", imageName, err)
return err
@@ -2202,25 +2444,412 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error {
return nil
}
+func findActiveSwarmNodes(dockercli *dockerclient.Client) (int64, error) {
+ ctx := context.Background()
+ nodes, err := dockercli.NodeList(ctx, types.NodeListOptions{})
+ if err != nil {
+ return 0, err
+ }
+
+ nodeCount := int64(0)
+ for _, node := range nodes {
+ //log.Printf("ID: %s - %#v", node.ID, node.Status.State)
+ if node.Status.State == "ready" {
+ nodeCount += 1
+ }
+ }
+
+ return nodeCount, nil
+
+ /*
+ containers, err := dockercli.ContainerList(ctx, types.ContainerListOptions{
+ All: true,
+ })
+ */
+}
+
+
+// Runs data discovery
+
+func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, action *shuffle.Action, workflowExecution *shuffle.WorkflowExecution) error {
+ parsedRequest := shuffle.OrborusExecutionRequest{
+ Cleanup: cleanupEnv,
+ ExecutionId: workflowExecution.ExecutionId,
+ Authorization: workflowExecution.Authorization,
+ EnvironmentName: os.Getenv("ENVIRONMENT_NAME"),
+ Timezone: os.Getenv("TZ"),
+ HTTPProxy: os.Getenv("HTTP_PROXY"),
+ HTTPSProxy: os.Getenv("HTTPS_PROXY"),
+ ShufflePassProxyToApp: os.Getenv("SHUFFLE_PASS_APP_PROXY"),
+ Url: baseUrl,
+ BaseUrl: baseUrl,
+ Action: *action,
+ FullExecution: *workflowExecution,
+ }
+ // Sometimes makes it have the wrong data due to timing
+
+ // Specific for subflow to ensure worker matches the backend correctly
+
+ parsedBaseurl := incomingUrl
+ if strings.Count(baseUrl, ":") >= 2 {
+ baseUrlSplit := strings.Split(baseUrl, ":")
+ if len(baseUrlSplit) >= 3 {
+ parsedBaseurl = strings.Join(baseUrlSplit[0:2], ":")
+ //parsedRequest.BaseUrl = fmt.Sprintf("%s:33333", parsedBaseurl)
+ }
+ }
+
+ if len(parsedRequest.Url) == 0 {
+ // Fixed callback url to the worker itself
+ if strings.Count(parsedBaseurl, ":") >= 2 {
+ parsedRequest.Url = parsedBaseurl
+ } else {
+ // Callback to worker
+ parsedRequest.Url = fmt.Sprintf("%s:%d", parsedBaseurl, baseport)
+
+ //parsedRequest.Url
+ }
+
+ //log.Printf("[DEBUG][%s] Should add a baseurl for the app to get back to: %s", workflowExecution.ExecutionId, parsedRequest.Url)
+ }
+
+ // Swapping because this was confusing during dev
+ // No real reason, just variable names
+ tmp := parsedRequest.Url
+ parsedRequest.Url = parsedRequest.BaseUrl
+ parsedRequest.BaseUrl = tmp
+
+ // Run with proper hostname, but set to shuffle-worker to avoid specific host target.
+ // This means running with VIP instead.
+ if len(hostname) > 0 {
+ parsedRequest.BaseUrl = fmt.Sprintf("http://%s:%d", hostname, baseport)
+ //parsedRequest.BaseUrl = fmt.Sprintf("http://shuffle-workers:%d", baseport)
+ //log.Printf("[DEBUG][%s] Changing hostname to local hostname in Docker network for WORKER URL: %s", workflowExecution.ExecutionId, parsedRequest.BaseUrl)
+
+ if parsedRequest.Action.AppName == "shuffle-subflow" || parsedRequest.Action.AppName == "shuffle-subflow-v2" || parsedRequest.Action.AppName == "User Input" {
+ parsedRequest.BaseUrl = fmt.Sprintf("http://%s:%d", hostname, baseport)
+ //parsedRequest.Url = parsedRequest.BaseUrl
+ }
+ }
+
+ // Making sure to get the LATEST execution data
+ // This is due to cache timing issues
+ exec, err := shuffle.GetWorkflowExecution(ctx, workflowExecution.ExecutionId)
+ if err == nil && len(exec.ExecutionId) > 0 {
+ parsedRequest.FullExecution = *exec
+ }
+
+ data, err := json.Marshal(parsedRequest)
+ if err != nil {
+ log.Printf("[ERROR] Failed marshalling worker request: %s", err)
+ return err
+ }
+
+ streamUrl := fmt.Sprintf("http://%s:%d/api/v1/run", appName, port)
+ log.Printf("[DEBUG][%s] Worker URL: %s, Backend URL: %s, Target App: %s", workflowExecution.ExecutionId, parsedRequest.BaseUrl, parsedRequest.Url, streamUrl)
+ req, err := http.NewRequest(
+ "POST",
+ streamUrl,
+ bytes.NewBuffer([]byte(data)),
+ )
+
+ // Checking as LATE as possible, ensuring we don't rerun what's already ran
+ //ctx = context.Background()
+ newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
+ _, err = shuffle.GetCache(ctx, newExecId)
+ if err == nil {
+ log.Printf("[DEBUG] Result for %s already found (PRE REQUEST) - returning", newExecId)
+ return nil
+ }
+
+ cacheData := []byte("1")
+ err = shuffle.SetCache(ctx, newExecId, cacheData, 30)
+ if err != nil {
+ log.Printf("[WARNING] Failed setting cache for action %s: %s", newExecId, err)
+ } else {
+ log.Printf("[DEBUG][%s] Adding %s to cache (%#v)", workflowExecution.ExecutionId, newExecId, action.Name)
+ }
+
+ // FIXME: Add 5 tries
+
+ newresp, err := topClient.Do(req)
+ if err != nil {
+ // Another timeout issue here somewhere
+ // context deadline
+ if strings.Contains(fmt.Sprintf("%s", err), "context deadline exceeded") || strings.Contains(fmt.Sprintf("%s", err), "Client.Timeout exceeded") {
+ return nil
+ }
+
+ if strings.Contains(fmt.Sprintf("%s", err), "timeout awaiting response") {
+ return nil
+ }
+
+ newerr := fmt.Sprintf("%s", err)
+ if strings.Contains(newerr, "connection refused") || strings.Contains(newerr, "no such host") {
+ newerr = fmt.Sprintf("Failed connecting to app %s. Is the Docker image available?", appName)
+ } else {
+ // escape quotes and newlines
+ newerr = strings.ReplaceAll(strings.ReplaceAll(newerr, "\"", "\\\""), "\n", "\\n")
+ }
+
+ log.Printf("[ERROR][%s] Error running app run request: %s", workflowExecution.ExecutionId, err)
+ actionResult := shuffle.ActionResult{
+ Action: *action,
+ ExecutionId: workflowExecution.ExecutionId,
+ Authorization: workflowExecution.Authorization,
+ Result: fmt.Sprintf(`{"success": false, "reason": "Failed to connect to app %s in swarm. Restart Orborus if this is recurring, or contact support@shuffler.io.", "details": "%s"}`, streamUrl, newerr),
+ StartedAt: int64(time.Now().Unix()),
+ CompletedAt: int64(time.Now().Unix()),
+ Status: "FAILURE",
+ }
+
+ // If this happens - send failure signal to stop the workflow?
+ sendSelfRequest(actionResult)
+ return err
+ }
+
+ defer newresp.Body.Close()
+ body, err := ioutil.ReadAll(newresp.Body)
+ if err != nil {
+ log.Printf("[ERROR] Failed reading app request body body: %s", err)
+ return err
+ } else {
+ log.Printf("[DEBUG][%s] NEWRESP (from app): %s", workflowExecution.ExecutionId, string(body))
+ }
+
+ return nil
+}
+
+// Function to auto-deploy certain apps if "run" is set
+// Has some issues with loading when running multiple workers and such.
+func baseDeploy() {
+
+ cli, err := dockerclient.NewEnvClient()
+ if err != nil {
+ log.Printf("[ERROR] Unable to create docker client (3): %s", err)
+ return
+ }
+
+ for key, value := range autoDeploy {
+ newNameSplit := strings.Split(key, ":")
+
+ action := shuffle.Action{
+ AppName: newNameSplit[0],
+ AppVersion: newNameSplit[1],
+ ID: "TBD",
+ }
+
+ workflowExecution := shuffle.WorkflowExecution{
+ ExecutionId: "TBD",
+ }
+
+ appname := action.AppName
+ appversion := action.AppVersion
+ appname = strings.Replace(appname, ".", "-", -1)
+ appversion = strings.Replace(appversion, ".", "-", -1)
+
+ env := []string{
+ fmt.Sprintf("EXECUTIONID=%s", workflowExecution.ExecutionId),
+ fmt.Sprintf("AUTHORIZATION=%s", workflowExecution.Authorization),
+ fmt.Sprintf("CALLBACK_URL=%s", baseUrl),
+ fmt.Sprintf("BASE_URL=%s", appCallbackUrl),
+ fmt.Sprintf("TZ=%s", timezone),
+ fmt.Sprintf("SHUFFLE_LOGS_DISABLED=%s", logsDisabled),
+ }
+
+ if strings.ToLower(os.Getenv("SHUFFLE_PASS_APP_PROXY")) == "true" {
+ //log.Printf("APPENDING PROXY TO THE APP!")
+ env = append(env, fmt.Sprintf("HTTP_PROXY=%s", os.Getenv("HTTP_PROXY")))
+ env = append(env, fmt.Sprintf("HTTPS_PROXY=%s", os.Getenv("HTTPS_PROXY")))
+ env = append(env, fmt.Sprintf("NO_PROXY=%s", os.Getenv("NO_PROXY")))
+ }
+
+ if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 {
+ log.Printf("[DEBUG] Setting SHUFFLE_APP_SDK_TIMEOUT to %s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT"))
+ env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
+ }
+
+ identifier := fmt.Sprintf("%s_%s_%s_%s", appname, appversion, action.ID, workflowExecution.ExecutionId)
+ if strings.Contains(identifier, " ") {
+ identifier = strings.ReplaceAll(identifier, " ", "-")
+ }
+
+ //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)
+ err = deployApp(cli, value, identifier, env, workflowExecution, action)
+ _ = err
+ //err := deployApp(cli, value, identifier, env, workflowExecution, action)
+ //if err != nil {
+ // log.Printf("[DEBUG] Failed deploying app %s: %s", value, err)
+ //}
+ }
+
+ appsInitialized = true
+}
+
+func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExecution shuffle.WorkflowExecution, firstRequest bool, environments []string) ([]string, error) {
+ // Because of this, it always has updated data.
+ // Removed request requirement from app_sdk
+ newresp, err := topClient.Do(req)
+ if err != nil {
+ log.Printf("[ERROR] Failed request: %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return environments, err
+ }
+
+ defer newresp.Body.Close()
+ body, err := ioutil.ReadAll(newresp.Body)
+ if err != nil {
+ log.Printf("[ERROR] Failed reading body: %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return environments, err
+ }
+
+ if newresp.StatusCode != 200 {
+ log.Printf("[ERROR] %sStatusCode (1): %d", string(body), newresp.StatusCode)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode) )
+ }
+
+ err = json.Unmarshal(body, &workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err)
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return environments, err
+ }
+
+ if firstRequest {
+ firstRequest = false
+
+ ctx := context.Background()
+ cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId)
+ execData, err := json.Marshal(workflowExecution)
+ if err != nil {
+ log.Printf("[ERROR][%s] Failed marshalling execution during set (3): %s", workflowExecution.ExecutionId, err)
+ } else {
+ err = shuffle.SetCache(ctx, cacheKey, execData, 30)
+ if err != nil {
+ log.Printf("[ERROR][%s] Failed adding to cache during setexecution (3): %s", workflowExecution.ExecutionId, err)
+ }
+ }
+
+ 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)
+ }
+ }
+
+ // Checks if a subflow is child of the startnode, as sub-subflows aren't working properly yet
+ childNodes := shuffle.FindChildNodes(workflowExecution, workflowExecution.Start, []string{}, []string{})
+ log.Printf("[DEBUG] Looking for subflow in %#v to check execution pattern as child of %s", childNodes, workflowExecution.Start)
+ subflowFound := false
+ for _, childNode := range childNodes {
+ for _, trigger := range workflowExecution.Workflow.Triggers {
+ if trigger.ID != childNode {
+ continue
+ }
+
+ if trigger.AppName == "Shuffle Workflow" {
+ subflowFound = true
+ break
+ }
+ }
+
+ if subflowFound {
+ break
+ }
+ }
+
+ log.Printf("[DEBUG] Environments: %s. Source: %s. 1 env = webserver, 0 or >1 = default. Subflow exists: %#v", environments, workflowExecution.ExecutionSource, subflowFound)
+ if len(environments) == 1 && workflowExecution.ExecutionSource != "default" && !subflowFound {
+ log.Printf("[DEBUG] Running OPTIMIZED execution (not manual)")
+ listener := webserverSetup(workflowExecution)
+ err := executionInit(workflowExecution)
+ if err != nil {
+ log.Printf("[DEBUG] 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)
+ //log.Printf("Before wait")
+ //wg := sync.WaitGroup{}
+ //wg.Add(1)
+ //wg.Wait()
+ } else {
+ log.Printf("[DEBUG] Running NON-OPTIMIZED execution for type %s with %d environment(s). This only happens when ran manually OR when running with subflows. Status: %s", workflowExecution.ExecutionSource, len(environments), workflowExecution.Status)
+ err := executionInit(workflowExecution)
+ if err != nil {
+ log.Printf("[DEBUG] Workflow setup failed: %s", workflowExecution.ExecutionId, err)
+ shutdown(workflowExecution, "", "", true)
+ }
+
+ // Trying to make worker into microservice~ :)
+ }
+ }
+
+ if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
+ log.Printf("[DEBUG] Workflow %s is finished. Exiting worker.", workflowExecution.ExecutionId)
+ log.Printf("[DEBUG] Shutting down (31)")
+ shutdown(workflowExecution, "", "", true)
+ }
+
+ if workflowExecution.Status == "EXECUTING" || workflowExecution.Status == "RUNNING" {
+ //log.Printf("Status: %s", workflowExecution.Status)
+ err = handleDefaultExecution(client, req, workflowExecution)
+ if err != nil {
+ log.Printf("[DEBUG] Workflow %s is finished: %s", workflowExecution.ExecutionId, err)
+ log.Printf("[DEBUG] Shutting down (32)")
+ shutdown(workflowExecution, "", "", true)
+ }
+ } else {
+ log.Printf("[DEBUG] Workflow %s has status %s. Exiting worker (if WAITING, rerun will happen).", workflowExecution.ExecutionId, workflowExecution.Status)
+ log.Printf("[DEBUG] Shutting down (33)")
+ shutdown(workflowExecution, workflowExecution.Workflow.ID, "", true)
+ }
+
+ time.Sleep(time.Duration(sleepTime) * time.Second)
+ return environments, nil
+}
+
// Initial loop etc
func main() {
// Elasticsearch necessary to ensure we'ren ot running with Datastore configurations for minimal/maximal data sizes
- _, err := shuffle.RunInit(datastore.Client{}, storage.Client{}, "", "worker", true, "elasticsearch")
+ // Recursive import kind of :)
+ _, err := shuffle.RunInit(*shuffle.GetDatastore(), *shuffle.GetStorage(), "", "worker", true, "elasticsearch")
if err != nil {
- log.Printf("[ERROR] Failed to run worker init: %s", err)
+ if !strings.Contains(fmt.Sprintf("%s", err), "no such host") {
+ 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
+ sleepTime = 5
client := shuffle.GetExternalClient(baseUrl)
if timezone == "" {
timezone = "Europe/Amsterdam"
}
- log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, os.Getenv("SHUFFLE_SWARM_CONFIG"))
+ topClient = client
+ swarmConfig := os.Getenv("SHUFFLE_SWARM_CONFIG")
+ log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, swarmConfig)
+
authorization := ""
executionId := ""
@@ -2269,148 +2898,13 @@ func main() {
}
topClient = client
-
firstRequest := true
environments := []string{}
for {
- // Because of this, it always has updated data.
- // Removed request requirement from app_sdk
- newresp, err := client.Do(req)
+ environments, err = getStreamResultsWrapper(client, req, workflowExecution, firstRequest, environments)
if err != nil {
- log.Printf("[ERROR] Failed request: %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
+ log.Printf("[ERROR] Failed getting stream results: %s", err)
}
-
- defer newresp.Body.Close()
- body, err := ioutil.ReadAll(newresp.Body)
- if err != nil {
- log.Printf("[ERROR] Failed reading body: %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- if newresp.StatusCode != 200 {
- log.Printf("[ERROR] %s\nStatusCode (1): %d", string(body), newresp.StatusCode)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- err = json.Unmarshal(body, &workflowExecution)
- if err != nil {
- log.Printf("[ERROR] Failed workflowExecution unmarshal: %s", err)
- time.Sleep(time.Duration(sleepTime) * time.Second)
- continue
- }
-
- if firstRequest {
- firstRequest = false
- //workflowExecution.StartedAt = int64(time.Now().Unix())
-
- ctx := context.Background()
- cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId)
- execData, err := json.Marshal(workflowExecution)
- if err != nil {
- log.Printf("[ERROR][%s] Failed marshalling execution during set (3): %s", workflowExecution.ExecutionId, err)
- } else {
- err = shuffle.SetCache(ctx, cacheKey, execData, 30)
- if err != nil {
- log.Printf("[ERROR][%s] Failed adding to cache during setexecution (3): %s", workflowExecution.ExecutionId, err)
- }
- }
-
- //requestCache = cache.New(60*time.Minute, 120*time.Minute)
- //requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration)
-
- 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)
- }
- }
-
- // Checks if a subflow is child of the startnode, as sub-subflows aren't working properly yet
- childNodes := shuffle.FindChildNodes(workflowExecution, workflowExecution.Start, []string{}, []string{})
- log.Printf("[DEBUG] Looking for subflow in %#v to check execution pattern as child of %s", childNodes, workflowExecution.Start)
- subflowFound := false
- for _, childNode := range childNodes {
- for _, trigger := range workflowExecution.Workflow.Triggers {
- if trigger.ID != childNode {
- continue
- }
-
- if trigger.AppName == "Shuffle Workflow" {
- subflowFound = true
- break
- }
- }
-
- if subflowFound {
- break
- }
- }
-
- log.Printf("\n\nEnvironments: %s. Source: %s. 1 env = webserver, 0 or >1 = default. Subflow exists: %#v\n\n", environments, workflowExecution.ExecutionSource, subflowFound)
- if len(environments) == 1 && workflowExecution.ExecutionSource != "default" && !subflowFound {
- log.Printf("\n\n[DEBUG] Running OPTIMIZED execution (not manual)\n\n")
- listener := webserverSetup(workflowExecution)
- err := executionInit(workflowExecution)
- if err != nil {
- log.Printf("[DEBUG] 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)
- //log.Printf("Before wait")
- //wg := sync.WaitGroup{}
- //wg.Add(1)
- //wg.Wait()
- } else {
- log.Printf("\n\n[DEBUG] 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("[DEBUG] Workflow setup failed: %s", workflowExecution.ExecutionId, err)
- shutdown(workflowExecution, "", "", true)
- }
-
- // Trying to make worker into microservice~ :)
- }
- }
-
- if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "SUCCESS" {
- log.Printf("[DEBUG] Workflow %s is finished. Exiting worker.", workflowExecution.ExecutionId)
- log.Printf("[DEBUG] Shutting down (31)")
- shutdown(workflowExecution, "", "", true)
- }
-
- if workflowExecution.Status == "EXECUTING" || workflowExecution.Status == "RUNNING" {
- //log.Printf("Status: %s", workflowExecution.Status)
- err = handleDefaultExecution(client, req, workflowExecution)
- if err != nil {
- log.Printf("[DEBUG] Workflow %s is finished: %s", workflowExecution.ExecutionId, err)
- log.Printf("[DEBUG] Shutting down (32)")
- shutdown(workflowExecution, "", "", true)
- }
- } else {
- log.Printf("[DEBUG] Workflow %s has status %s. Exiting worker.", workflowExecution.ExecutionId, workflowExecution.Status)
- log.Printf("[DEBUG] Shutting down (33)")
- shutdown(workflowExecution, workflowExecution.Workflow.ID, "", true)
- }
-
- time.Sleep(time.Duration(sleepTime) * time.Second)
}
}
@@ -2447,6 +2941,7 @@ func checkUnfinished(resp http.ResponseWriter, request *http.Request, execReques
}
func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
+ defer request.Body.Close()
body, err := ioutil.ReadAll(request.Body)
if err != nil {
log.Printf("[WARNING] Failed reading body for stream result queue")
@@ -2455,8 +2950,6 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
return
}
- defer request.Body.Close()
-
//log.Printf("[DEBUG] In run execution with body length %d", len(body))
var execRequest shuffle.OrborusExecutionRequest
err = json.Unmarshal(body, &execRequest)
@@ -2511,10 +3004,10 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
os.Setenv("AUTHORIZATION", execRequest.Authorization)
}
- 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)
+ topClient = shuffle.GetExternalClient(streamResultUrl)
req, err := http.NewRequest(
"POST",
@@ -2530,6 +3023,7 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
return
}
+ defer newresp.Body.Close()
body, err = ioutil.ReadAll(newresp.Body)
if err != nil {
log.Printf("[ERROR] Failed reading body (2): %s", err)
@@ -2560,6 +3054,7 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
}
ctx := context.Background()
+ //err = shuffle.SetWorkflowExecution(ctx, workflowExecution, true)
err = setWorkflowExecution(ctx, workflowExecution, true)
if err != nil {
log.Printf("[ERROR] Failed initializing execution saving for %s: %s", workflowExecution.ExecutionId, err)
@@ -2601,14 +3096,12 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
if err != nil {
log.Printf("[ERROR][%s] Failed marshalling execution during set (3): %s", workflowExecution.ExecutionId, err)
} else {
- err = shuffle.SetCache(ctx, cacheKey, execData, 30)
+ err = shuffle.SetCache(ctx, cacheKey, execData, 31)
if err != nil {
log.Printf("[ERROR][%s] Failed adding to cache during setexecution (3): %s", workflowExecution.ExecutionId, err)
}
}
- //requestCache.Set(cacheKey, &workflowExecution, cache.DefaultExpiration)
-
err = executionInit(workflowExecution)
if err != nil {
log.Printf("[DEBUG][%s] Shutting down (30) - Workflow setup failed: %s", workflowExecution.ExecutionId, workflowExecution.ExecutionId, err)
@@ -2618,16 +3111,95 @@ func handleRunExecution(resp http.ResponseWriter, request *http.Request) {
//shutdown(workflowExecution, "", "", true)
}
- //go handleExecutionResult(workflowExecution)
handleExecutionResult(workflowExecution)
resp.WriteHeader(200)
resp.Write([]byte(fmt.Sprintf(`{"success": true}`)))
}
+func handleDownloadImage(resp http.ResponseWriter, request *http.Request) {
+ // Read the request body
+ defer request.Body.Close()
+ bodyBytes, err := ioutil.ReadAll(request.Body)
+ if err != nil {
+ log.Printf("[ERROR] Failed reading body for stream result queue. Error: %s", err)
+ resp.WriteHeader(401)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
+ return
+ }
+
+ // get images from request
+ image := &ImageDownloadBody{}
+ err = json.Unmarshal(bodyBytes, image)
+ if err != nil {
+ log.Printf("[ERROR] Error in unmarshalling body: %s", err)
+ resp.WriteHeader(401)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
+ return
+ }
+
+ client, err := dockerclient.NewEnvClient()
+ if err != nil {
+ log.Printf("[ERROR] Unable to create docker client (4): %s", err)
+ resp.WriteHeader(401)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
+ return
+ }
+
+ // check if images are already downloaded
+ // Retrieve a list of Docker images
+ images, err := client.ImageList(context.Background(), types.ImageListOptions{})
+ if err != nil {
+ log.Printf("[ERROR] listing images: %s", err)
+ resp.WriteHeader(401)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
+ return
+ }
+
+ for _, img := range images {
+ for _, tag := range img.RepoTags {
+ splitTag := strings.Split(tag, ":")
+ baseTag := tag
+ if len(splitTag) > 1 {
+ baseTag = splitTag[1]
+ }
+
+ var possibleNames []string
+ possibleNames = append(possibleNames, fmt.Sprintf("frikky/shuffle:%s", baseTag))
+ possibleNames = append(possibleNames, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", baseTag))
+
+ if (arrayContains(possibleNames, image.Image)) {
+ log.Printf("[DEBUG] Image %s already downloaded that has been requested to download", image.Image)
+ resp.WriteHeader(200)
+ resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "image already present"}`)))
+ return
+ }
+ }
+ }
+
+ log.Printf("[INFO] Downloading image %s", image.Image)
+ downloadDockerImageBackend(&http.Client{Timeout: 60 * time.Second}, image.Image)
+
+ // return success
+ resp.WriteHeader(200)
+ resp.Write([]byte(fmt.Sprintf(`{"success": true, "status": "starting download"}`)))
+}
+
func runWebserver(listener net.Listener) {
r := mux.NewRouter()
r.HandleFunc("/api/v1/streams", handleWorkflowQueue).Methods("POST", "OPTIONS")
r.HandleFunc("/api/v1/streams/results", handleGetStreamResults).Methods("POST", "OPTIONS")
+ r.HandleFunc("/api/v1/execute", handleRunExecution).Methods("POST", "OPTIONS")
+ r.HandleFunc("/api/v1/run", handleRunExecution).Methods("POST", "OPTIONS")
+ r.HandleFunc("/api/v1/download", handleDownloadImage).Methods("POST", "OPTIONS")
+
+
+ if strings.ToLower(os.Getenv("SHUFFLE_DEBUG_MEMORY")) == "true" {
+ r.HandleFunc("/debug/pprof/", pprof.Index)
+ r.HandleFunc("/debug/pprof/heap", pprof.Handler("heap").ServeHTTP)
+ r.HandleFunc("/debug/pprof/profile", pprof.Profile)
+ r.HandleFunc("/debug/pprof/symbol", pprof.Symbol)
+ r.HandleFunc("/debug/pprof/trace", pprof.Trace)
+ }
//log.Fatal(http.ListenAndServe(port, nil))
//srv := http.Server{
@@ -2638,7 +3210,7 @@ func runWebserver(listener net.Listener) {
//log.Fatal(http.Serve(listener, nil))
- log.Printf("\n\n[DEBUG] NEW webserver setup\n\n")
+ log.Printf("[DEBUG] NEW webserver setup")
http.Handle("/", r)
srv := http.Server{
@@ -2651,7 +3223,7 @@ func runWebserver(listener net.Listener) {
err := srv.Serve(listener)
if err != nil {
- log.Printf("serveIssue: %#v", err)
+ log.Printf("[ERROR] Serve issue in worker: %#v", err)
}
log.Printf("[DEBUG] Do we see this?")
}