|
|
|
@@ -3,6 +3,7 @@ package main
|
|
|
|
|
import (
|
|
|
|
|
"github.com/shuffle/shuffle-shared"
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
"bytes"
|
|
|
|
|
"context"
|
|
|
|
|
"encoding/json"
|
|
|
|
@@ -21,8 +22,8 @@ import (
|
|
|
|
|
"time"
|
|
|
|
|
|
|
|
|
|
"github.com/docker/docker/api/types"
|
|
|
|
|
"github.com/docker/docker/api/types/container"
|
|
|
|
|
"github.com/docker/docker/api/types/filters"
|
|
|
|
|
"github.com/docker/docker/api/types/container"
|
|
|
|
|
"github.com/docker/docker/api/types/mount"
|
|
|
|
|
dockerclient "github.com/docker/docker/client"
|
|
|
|
|
// This is for automatic removal of certain code :)
|
|
|
|
@@ -34,10 +35,6 @@ import (
|
|
|
|
|
corev1 "k8s.io/api/core/v1"
|
|
|
|
|
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
|
|
|
|
"k8s.io/client-go/kubernetes"
|
|
|
|
|
"k8s.io/client-go/rest"
|
|
|
|
|
"k8s.io/client-go/tools/clientcmd"
|
|
|
|
|
"k8s.io/client-go/util/homedir"
|
|
|
|
|
"path/filepath"
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
// This is getting out of hand :)
|
|
|
|
@@ -56,7 +53,6 @@ var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
|
|
|
|
|
|
|
|
|
|
// var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME")
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// var baseimagename = "registry.hub.docker.com/frikky/shuffle"
|
|
|
|
|
var registryName = "registry.hub.docker.com"
|
|
|
|
|
var sleepTime = 2
|
|
|
|
@@ -81,7 +77,6 @@ var startAction string
|
|
|
|
|
//var allLogs map[string]string
|
|
|
|
|
//var containerIds []string
|
|
|
|
|
var downloadedImages []string
|
|
|
|
|
|
|
|
|
|
type ImageDownloadBody struct {
|
|
|
|
|
Image string `json:"image"`
|
|
|
|
|
}
|
|
|
|
@@ -93,6 +88,7 @@ type ImageRequest struct {
|
|
|
|
|
var finishedExecutions []string
|
|
|
|
|
var imagesDistributed []string
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Images to be autodeployed in the latest version of Shuffle.
|
|
|
|
|
var autoDeploy = map[string]string{
|
|
|
|
|
"http:1.4.0": "frikky/shuffle:http_1.4.0",
|
|
|
|
@@ -139,6 +135,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
return err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
handleExecutionResult(workflowExecution)
|
|
|
|
|
validated := shuffle.ValidateFinished(ctx, -1, workflowExecution)
|
|
|
|
|
if validated {
|
|
|
|
@@ -178,7 +175,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if len(subflowId) == 0 {
|
|
|
|
|
if len(subflowId) == 0 {
|
|
|
|
|
log.Printf("[DEBUG][%s] No waiting result found. Not polling", workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
for _, action := range workflowExecution.Workflow.Actions {
|
|
|
|
@@ -186,17 +183,19 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
workflowExecution.Workflow.Triggers = append(workflowExecution.Workflow.Triggers, shuffle.Trigger{
|
|
|
|
|
AppName: action.AppName,
|
|
|
|
|
Parameters: action.Parameters,
|
|
|
|
|
ID: action.ID,
|
|
|
|
|
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 {
|
|
|
|
@@ -215,9 +214,9 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
//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
|
|
|
|
|
found = true
|
|
|
|
|
break
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -236,20 +235,21 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
|
|
|
|
|
if len(subflowId) > 0 {
|
|
|
|
|
// Under rerun period timeout
|
|
|
|
|
timeComparison := 120
|
|
|
|
|
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)
|
|
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
@@ -271,6 +271,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
|
|
|
|
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)
|
|
|
|
@@ -310,7 +311,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
|
|
|
|
|
}
|
|
|
|
|
*/
|
|
|
|
|
} else {
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if len(reason) > 0 && len(nodeId) > 0 {
|
|
|
|
@@ -393,7 +394,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
clientset, err := getKubernetesClient()
|
|
|
|
|
clientset, _, err := shuffle.GetKubernetesClient()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Failed getting kubernetes: %s", err)
|
|
|
|
|
return err
|
|
|
|
@@ -499,7 +500,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|
|
|
|
|
|
|
|
|
if !strings.Contains(param.Value, "shuffle-backend") {
|
|
|
|
|
continue
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Automatic replacement as this is default
|
|
|
|
|
if len(os.Getenv("BASE_URL")) > 0 {
|
|
|
|
@@ -514,6 +515,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Max 10% CPU every second
|
|
|
|
|
//CPUShares: 128,
|
|
|
|
|
//CPUQuota: 10000,
|
|
|
|
@@ -540,7 +542,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|
|
|
|
|
|
|
|
|
// Get environment for certificates
|
|
|
|
|
volumeBinds := []string{}
|
|
|
|
|
volumeBindString := os.Getenv("SHUFFLE_VOLUME_BINDS")
|
|
|
|
|
volumeBindString:= os.Getenv("SHUFFLE_VOLUME_BINDS")
|
|
|
|
|
if len(volumeBindString) > 0 {
|
|
|
|
|
volumeBindSplit := strings.Split(volumeBindString, ",")
|
|
|
|
|
for _, volumeBind := range volumeBindSplit {
|
|
|
|
@@ -581,6 +583,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|
|
|
|
Env: env,
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Checking as late as possible, just in case.
|
|
|
|
|
newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
|
|
|
|
|
_, err := shuffle.GetCache(ctx, newExecId)
|
|
|
|
@@ -860,7 +863,7 @@ func askOtherWorkersToDownloadImage(image string) {
|
|
|
|
|
// Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD
|
|
|
|
|
if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" {
|
|
|
|
|
log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. NOT distributing images %s", image)
|
|
|
|
|
return
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if shuffle.ArrayContains(imagesDistributed, image) {
|
|
|
|
@@ -893,7 +896,7 @@ func askOtherWorkersToDownloadImage(image string) {
|
|
|
|
|
req, err := http.NewRequest(
|
|
|
|
|
"POST",
|
|
|
|
|
url,
|
|
|
|
|
bytes.NewBuffer(imageJSON),
|
|
|
|
|
bytes.NewBuffer(imageJSON),
|
|
|
|
|
)
|
|
|
|
|
|
|
|
|
|
if err != nil {
|
|
|
|
@@ -933,6 +936,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
dockercli, err := dockerclient.NewEnvClient()
|
|
|
|
@@ -1000,7 +1004,7 @@ 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("[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)
|
|
|
|
@@ -1086,9 +1090,10 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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")))
|
|
|
|
|
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
|
|
|
|
@@ -1116,6 +1121,8 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
fmt.Sprintf("%s:%s_%s", baseimagename, parsedAppname, action.AppVersion),
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// If cleanup is set, it should run for efficiency
|
|
|
|
|
pullOptions := types.ImagePullOptions{}
|
|
|
|
|
if cleanupEnv == "true" {
|
|
|
|
@@ -1378,7 +1385,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
|
|
|
|
log.Printf("[DEBUG][%s] Shutting down (17)", workflowExecution.ExecutionId)
|
|
|
|
|
if isKubernetes == "true" {
|
|
|
|
|
// log.Printf("workflow execution: %#v", workflowExecution)
|
|
|
|
|
clientset, err := getKubernetesClient()
|
|
|
|
|
clientset, _, err := shuffle.GetKubernetesClient()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("[ERROR] Error getting kubernetes client (1):", err)
|
|
|
|
|
os.Exit(1)
|
|
|
|
@@ -1587,7 +1594,7 @@ func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.Workflow
|
|
|
|
|
log.Printf("[DEBUG] Shutting down (20)")
|
|
|
|
|
if isKubernetes == "true" {
|
|
|
|
|
// log.Printf("workflow execution: %#v", workflowExecution)
|
|
|
|
|
clientset, err := getKubernetesClient()
|
|
|
|
|
clientset, _, err := shuffle.GetKubernetesClient()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("[ERROR] Error getting kubernetes client (2):", err)
|
|
|
|
|
os.Exit(1)
|
|
|
|
@@ -1618,6 +1625,7 @@ func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.Workflow
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
if workflowExecution.Status == "WAITING" && workflowExecution.ExecutionSource != "default" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
|
|
|
|
|
log.Printf("[INFO][%s] Workflow execution is waiting. Exiting worker, as backend will restart it.", workflowExecution.ExecutionId)
|
|
|
|
|
shutdown(workflowExecution, "", "", true)
|
|
|
|
@@ -1683,7 +1691,7 @@ func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffl
|
|
|
|
|
log.Printf("[DEBUG] Shutting down (20)")
|
|
|
|
|
if isKubernetes == "true" {
|
|
|
|
|
// log.Printf("workflow execution: %#v", workflowExecution)
|
|
|
|
|
clientset, err := getKubernetesClient()
|
|
|
|
|
clientset, _, err := shuffle.GetKubernetesClient()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("[ERROR] Error getting kubernetes client (2):", err)
|
|
|
|
|
os.Exit(1)
|
|
|
|
@@ -1700,7 +1708,7 @@ func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffl
|
|
|
|
|
log.Printf("[DEBUG] Shutting down (21)")
|
|
|
|
|
if isKubernetes == "true" {
|
|
|
|
|
// log.Printf("workflow execution: %#v", workflowExecution)
|
|
|
|
|
clientset, err := getKubernetesClient()
|
|
|
|
|
clientset, _, err := shuffle.GetKubernetesClient()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Println("[ERROR] Error getting kubernetes client (3):", err)
|
|
|
|
|
os.Exit(1)
|
|
|
|
@@ -1888,62 +1896,6 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar {
|
|
|
|
|
return envVars
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func getKubernetesClient() (*kubernetes.Clientset, error) {
|
|
|
|
|
|
|
|
|
|
// Gets the config content from Orborus.
|
|
|
|
|
kubeconfigContent := os.Getenv("KUBERNETES_CONFIG")
|
|
|
|
|
if len(kubeconfigContent) > 0 {
|
|
|
|
|
log.Printf("[INFO] Using KUBERNETES_CONFIG to set up Kubernetes client: %#v", os.Getenv("KUBERNETES_CONFIG"))
|
|
|
|
|
config, err := rest.InClusterConfig()
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Failed to create Kubernetes client from in-cluster config: %s", err)
|
|
|
|
|
} else {
|
|
|
|
|
// Replace client configuration with kubeconfig content
|
|
|
|
|
config, err = clientcmd.RESTConfigFromKubeConfig([]byte(kubeconfigContent))
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR] Failed to create Kubernetes client from KUBERNETES_CONFIG: %s", err)
|
|
|
|
|
} else {
|
|
|
|
|
// Create Kubernetes client
|
|
|
|
|
clientset, err := kubernetes.NewForConfig(config)
|
|
|
|
|
if err != nil {
|
|
|
|
|
return nil, err
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
return clientset, nil
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
// Fallback
|
|
|
|
|
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
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
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)
|
|
|
|
@@ -2050,9 +2002,10 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
|
|
|
|
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
|
|
|
|
|
// Force sending result
|
|
|
|
|
shutdownData, err := json.Marshal(workflowExecution)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR][%s] Failed marshalling execution (35): %s", workflowExecution.ExecutionId, err)
|
|
|
|
@@ -2128,11 +2081,11 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
|
|
|
|
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 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
|
|
|
|
|
}
|
|
|
|
|
*/
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
@@ -2165,6 +2118,8 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
func sendSelfRequest(actionResult shuffle.ActionResult) {
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
data, err := json.Marshal(actionResult)
|
|
|
|
|
if err != nil {
|
|
|
|
|
log.Printf("[ERROR][%s] Shutting down (24): Failed to unmarshal data for backend: %s", actionResult.ExecutionId, err)
|
|
|
|
@@ -2223,10 +2178,10 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
|
|
|
|
|
|
|
|
|
// 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
|
|
|
|
|
}
|
|
|
|
|
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
|
|
|
|
@@ -2237,7 +2192,6 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
finishedExecutions = append(finishedExecutions, workflowExecution.ExecutionId)
|
|
|
|
|
|
|
|
|
|
*/
|
|
|
|
|
|
|
|
|
|
streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl)
|
|
|
|
@@ -2294,7 +2248,6 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
|
|
|
|
|
|
|
|
|
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 {
|
|
|
|
|
if result.Status == "EXECUTING" || result.Status == "WAITING" {
|
|
|
|
@@ -2303,7 +2256,8 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
log.Printf("[DEBUG][%s] Should send full result to %s", workflowExecution.ExecutionId, baseUrl)
|
|
|
|
|
|
|
|
|
|
//data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization)
|
|
|
|
@@ -2387,7 +2341,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
// GetLocalIP returns the non loopback local IP of the host
|
|
|
|
|
func getLocalIP() string {
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
addrs, err := net.InterfaceAddrs()
|
|
|
|
|
if err != nil {
|
|
|
|
|
return ""
|
|
|
|
@@ -2432,7 +2386,8 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("[DEBUG] OLD HOSTNAME: %s", appCallbackUrl)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
port := listener.Addr().(*net.TCPAddr).Port
|
|
|
|
|
// Set the port environment variable
|
|
|
|
|
os.Setenv("WORKER_PORT", fmt.Sprintf("%d", port))
|
|
|
|
@@ -2447,22 +2402,21 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
|
|
|
|
|
func downloadDockerImageBackend(client *http.Client, imageName string) error {
|
|
|
|
|
// Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD
|
|
|
|
|
if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" {
|
|
|
|
|
//log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. Not downloading image %s", imageName)
|
|
|
|
|
log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. Not downloading image %s", imageName)
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
if arrayContains(downloadedImages, imageName) {
|
|
|
|
|
log.Printf("[DEBUG] Image %s already downloaded", imageName)
|
|
|
|
|
log.Printf("[DEBUG] Image %s already downloaded - not re-downloading", imageName)
|
|
|
|
|
return nil
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
log.Printf("[DEBUG] Trying to download image %s from backend %s as it doesn't exist. All images: %#v", imageName, baseUrl, downloadedImages)
|
|
|
|
|
|
|
|
|
|
downloadedImages = append(downloadedImages, imageName)
|
|
|
|
|
|
|
|
|
|
data := fmt.Sprintf(`{"name": "%s"}`, imageName)
|
|
|
|
|
dockerImgUrl := fmt.Sprintf("%s/api/v1/get_docker_image", baseUrl)
|
|
|
|
|
|
|
|
|
|
log.Printf("[DEBUG] Trying to download image %s from backend %s as it doesn't exist. Data sent: %#v, All images: %#v", imageName, baseUrl, data, downloadedImages)
|
|
|
|
|
|
|
|
|
|
req, err := http.NewRequest(
|
|
|
|
|
"POST",
|
|
|
|
@@ -2591,6 +2545,7 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error {
|
|
|
|
|
*/
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
// Runs data discovery
|
|
|
|
|
|
|
|
|
|
func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, action *shuffle.Action, workflowExecution *shuffle.WorkflowExecution) error {
|
|
|
|
@@ -2854,7 +2809,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe
|
|
|
|
|
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))
|
|
|
|
|
return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode) )
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
err = json.Unmarshal(body, &workflowExecution)
|
|
|
|
@@ -2937,6 +2892,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe
|
|
|
|
|
|
|
|
|
|
// Set environment variable
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
//log.Printf("Before wait")
|
|
|
|
|
//wg := sync.WaitGroup{}
|
|
|
|
|
//wg.Add(1)
|
|
|
|
@@ -3008,6 +2964,7 @@ func main() {
|
|
|
|
|
swarmConfig := os.Getenv("SHUFFLE_SWARM_CONFIG")
|
|
|
|
|
log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, swarmConfig)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
authorization := ""
|
|
|
|
|
executionId := ""
|
|
|
|
|
|
|
|
|
@@ -3312,6 +3269,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
return
|
|
|
|
|
}
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
for _, img := range images {
|
|
|
|
|
for _, tag := range img.RepoTags {
|
|
|
|
|
splitTag := strings.Split(tag, ":")
|
|
|
|
@@ -3324,7 +3282,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) {
|
|
|
|
|
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) {
|
|
|
|
|
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"}`)))
|
|
|
|
@@ -3349,6 +3307,7 @@ func runWebserver(listener net.Listener) {
|
|
|
|
|
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)
|
|
|
|
|