diff --git a/backend/go-app/go.sum b/backend/go-app/go.sum index 3297148f..7d974d2a 100644 --- a/backend/go-app/go.sum +++ b/backend/go-app/go.sum @@ -465,6 +465,8 @@ github.com/shuffle/shuffle-shared v0.6.61 h1:+9CCLeZLiAVDgNRTkZxnIgz+FZ7UrEHez2B github.com/shuffle/shuffle-shared v0.6.61/go.mod h1:RAJiSFjmuKmijKTbbEf9A6Ojb+3/te7g71lED7JjPus= github.com/shuffle/shuffle-shared v0.6.62 h1:NWVjVbnNpm6osDHQS2NGqwAgWk1Fk3s5RKUO3P5n4js= github.com/shuffle/shuffle-shared v0.6.62/go.mod h1:RAJiSFjmuKmijKTbbEf9A6Ojb+3/te7g71lED7JjPus= +github.com/shuffle/shuffle-shared v0.6.63 h1:eNQMpVhe/mAMxl61W9Wj6/Z4PrtPeEnbjvMtDdT1mqw= +github.com/shuffle/shuffle-shared v0.6.63/go.mod h1:RAJiSFjmuKmijKTbbEf9A6Ojb+3/te7g71lED7JjPus= github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0= github.com/sirupsen/logrus v1.9.0/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= diff --git a/functions/onprem/worker/go.mod b/functions/onprem/worker/go.mod index 38bf36bb..700a01dc 100644 --- a/functions/onprem/worker/go.mod +++ b/functions/onprem/worker/go.mod @@ -7,9 +7,9 @@ require ( github.com/gorilla/mux v1.8.1 github.com/satori/go.uuid v1.2.0 github.com/shuffle/shuffle-shared v0.6.63 - k8s.io/api v0.30.0 - k8s.io/apimachinery v0.30.0 - k8s.io/client-go v0.30.0 + k8s.io/api v0.30.2 + k8s.io/apimachinery v0.30.2 + k8s.io/client-go v0.30.2 ) require ( @@ -27,7 +27,7 @@ require ( github.com/algolia/algoliasearch-client-go/v3 v3.18.1 // indirect github.com/bradfitz/gomemcache v0.0.0-20230905024940-24af94b03874 // indirect github.com/bradfitz/slice v0.0.0-20180809154707-2b758aa73013 // indirect - github.com/cloudflare/circl v1.3.3 // indirect + github.com/cloudflare/circl v1.3.7 // indirect github.com/containerd/log v0.1.0 // indirect github.com/cyphar/filepath-securejoin v0.2.4 // indirect github.com/davecgh/go-spew v1.1.1 // indirect @@ -38,7 +38,7 @@ require ( github.com/emirpasic/gods v1.18.1 // indirect github.com/felixge/httpsnoop v1.0.4 // indirect github.com/frikky/kin-openapi v0.41.0 // indirect - github.com/frikky/schemaless v0.0.11 // indirect + github.com/frikky/schemaless v0.0.13 // indirect github.com/ghodss/yaml v1.0.0 // indirect github.com/go-git/gcfg v1.5.1-0.20230307220236-3a3c6141e376 // indirect github.com/go-git/go-billy/v5 v5.5.0 // indirect @@ -79,6 +79,8 @@ require ( github.com/pjbgf/sha1cd v0.3.0 // indirect github.com/pkg/errors v0.9.1 // indirect github.com/sashabaranov/go-openai v1.19.2 // indirect + github.com/sendgrid/rest v2.6.9+incompatible // indirect + github.com/sendgrid/sendgrid-go v3.14.0+incompatible // indirect github.com/sergi/go-diff v1.1.0 // indirect github.com/skeema/knownhosts v1.2.1 // indirect github.com/skip2/go-qrcode v0.0.0-20200617195104-da1b6568686e // indirect diff --git a/functions/onprem/worker/go.sum b/functions/onprem/worker/go.sum index 55e3a6c6..602a53a4 100644 --- a/functions/onprem/worker/go.sum +++ b/functions/onprem/worker/go.sum @@ -97,6 +97,8 @@ github.com/chzyer/test v0.0.0-20180213035817-a1ea475d72b1/go.mod h1:Q3SI9o4m/ZMn github.com/client9/misspell v0.3.4/go.mod h1:qj6jICC3Q7zFZvVWo7KLAzC3yx5G7kyvSDkc90ppPyw= github.com/cloudflare/circl v1.3.3 h1:fE/Qz0QdIGqeWfnwq0RE0R7MI51s0M2E4Ga9kq5AEMs= github.com/cloudflare/circl v1.3.3/go.mod h1:5XYMA4rFBvNIrhs50XuiBJ15vF2pZn4nnUKZrLbUZFA= +github.com/cloudflare/circl v1.3.7 h1:qlCDlTPz2n9fu58M0Nh1J/JzcFpfgkFHHX3O35r5vcU= +github.com/cloudflare/circl v1.3.7/go.mod h1:sRTcRWXGLrKw6yIGJ+l7amYJFfAXbZG0kBSc8r4zxgA= github.com/cncf/udpa/go v0.0.0-20191209042840-269d4d468f6f/go.mod h1:M8M6+tZqaGXZJjfX53e64911xZQV5JYwmTeXPW+k8Sc= github.com/cncf/udpa/go v0.0.0-20200629203442-efcf912fb354/go.mod h1:WmhPx2Nbnhtbo57+VJT5O0JRkEi1Wbu0z5j0R8u5Hbk= github.com/cncf/xds/go v0.0.0-20231128003011-0fa0005c9caa h1:jQCWAUqqlij9Pgj2i/PB79y4KOPYVyFYdROxgaCwdTQ= @@ -139,6 +141,8 @@ github.com/frikky/schemaless v0.0.9 h1:RzNLPkJq5c4nlm5iLiTndFcbeQxdMGJIj266wSGt2 github.com/frikky/schemaless v0.0.9/go.mod h1:mooDxY+D6weHjhKvjy3+IE9S7P4g4cpNnidkdRv/cHQ= github.com/frikky/schemaless v0.0.11 h1:c4r6CJX30XI+SoJdT9RlUd9qYSQlx6hvwGRtsypu+uM= github.com/frikky/schemaless v0.0.11/go.mod h1:mooDxY+D6weHjhKvjy3+IE9S7P4g4cpNnidkdRv/cHQ= +github.com/frikky/schemaless v0.0.13 h1:ARiN9V7wr2VZXAr9JK5wvTbyPgpGrgeiL1VhR5MlgaQ= +github.com/frikky/schemaless v0.0.13/go.mod h1:mooDxY+D6weHjhKvjy3+IE9S7P4g4cpNnidkdRv/cHQ= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ= github.com/ghodss/yaml v1.0.0 h1:wQHKEahhL6wmXdzwWG11gIVCkOv05bNOh+Rxn0yngAk= @@ -387,12 +391,18 @@ github.com/sashabaranov/go-openai v1.19.2 h1:+dkuCADSnwXV02YVJkdphY8XD9AyHLUWwk6 github.com/sashabaranov/go-openai v1.19.2/go.mod h1:lj5b/K+zjTSFxVLijLSTDZuP7adOgerWeFyZLUhAKRg= github.com/satori/go.uuid v1.2.0 h1:0uYX9dsZ2yD7q2RtLRtPSdGDWzjeM3TbMJP9utgA0ww= github.com/satori/go.uuid v1.2.0/go.mod h1:dA0hQrYB0VpLJoorglMZABFdXlWrHn1NEOzdhQKdks0= +github.com/sendgrid/rest v2.6.9+incompatible h1:1EyIcsNdn9KIisLW50MKwmSRSK+ekueiEMJ7NEoxJo0= +github.com/sendgrid/rest v2.6.9+incompatible/go.mod h1:kXX7q3jZtJXK5c5qK83bSGMdV6tsOE70KbHoqJls4lE= +github.com/sendgrid/sendgrid-go v3.14.0+incompatible h1:KDSasSTktAqMJCYClHVE94Fcif2i7P7wzISv1sU6DUA= +github.com/sendgrid/sendgrid-go v3.14.0+incompatible/go.mod h1:QRQt+LX/NmgVEvmdRw0VT/QgUn499+iza2FnDca9fg8= github.com/sergi/go-diff v1.1.0 h1:we8PVUC3FE2uYfodKH/nBHMSetSfHDR6scGdBi+erh0= github.com/sergi/go-diff v1.1.0/go.mod h1:STckp+ISIX8hZLjrqAeVduY0gWCT9IjLuqbuNXdaHfM= github.com/shuffle/shuffle-shared v0.6.16 h1:dQBDRmb2Wgl3pEuewqjDvN6v6nUKr+1EvGSEja9zG6s= github.com/shuffle/shuffle-shared v0.6.16/go.mod h1:HhQTn7xZZ69ZTc4EptO9OeNmgbKDyGlWAhFkUFUAHSA= github.com/shuffle/shuffle-shared v0.6.27 h1:q4qZD6bGZFIvZ5Y10unGr3N3rZ7OryWyvvaGgANZJZU= github.com/shuffle/shuffle-shared v0.6.27/go.mod h1:rWkh1eWdIx7OqQzJ1+JzF3Hck1X/Ty1WkUtjLrp+CU4= +github.com/shuffle/shuffle-shared v0.6.63 h1:eNQMpVhe/mAMxl61W9Wj6/Z4PrtPeEnbjvMtDdT1mqw= +github.com/shuffle/shuffle-shared v0.6.63/go.mod h1:RAJiSFjmuKmijKTbbEf9A6Ojb+3/te7g71lED7JjPus= github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0= github.com/sirupsen/logrus v1.9.0/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= @@ -906,10 +916,16 @@ honnef.co/go/tools v0.0.1-2020.1.3/go.mod h1:X/FiERA/W4tHapMX5mGpAtMSVEeEUOyHaw9 honnef.co/go/tools v0.0.1-2020.1.4/go.mod h1:X/FiERA/W4tHapMX5mGpAtMSVEeEUOyHaw9vFzvIQ3k= k8s.io/api v0.30.0 h1:siWhRq7cNjy2iHssOB9SCGNCl2spiF1dO3dABqZ8niA= k8s.io/api v0.30.0/go.mod h1:OPlaYhoHs8EQ1ql0R/TsUgaRPhpKNxIMrKQfWUp8QSE= +k8s.io/api v0.30.2 h1:+ZhRj+28QT4UOH+BKznu4CBgPWgkXO7XAvMcMl0qKvI= +k8s.io/api v0.30.2/go.mod h1:ULg5g9JvOev2dG0u2hig4Z7tQ2hHIuS+m8MNZ+X6EmI= k8s.io/apimachinery v0.30.0 h1:qxVPsyDM5XS96NIh9Oj6LavoVFYff/Pon9cZeDIkHHA= k8s.io/apimachinery v0.30.0/go.mod h1:iexa2somDaxdnj7bha06bhb43Zpa6eWH8N8dbqVjTUc= +k8s.io/apimachinery v0.30.2 h1:fEMcnBj6qkzzPGSVsAZtQThU62SmQ4ZymlXRC5yFSCg= +k8s.io/apimachinery v0.30.2/go.mod h1:iexa2somDaxdnj7bha06bhb43Zpa6eWH8N8dbqVjTUc= k8s.io/client-go v0.30.0 h1:sB1AGGlhY/o7KCyCEQ0bPWzYDL0pwOZO4vAtTSh/gJQ= k8s.io/client-go v0.30.0/go.mod h1:g7li5O5256qe6TYdAMyX/otJqMhIiGgTapdLchhmOaY= +k8s.io/client-go v0.30.2 h1:sBIVJdojUNPDU/jObC+18tXWcTJVcwyqS9diGdWHk50= +k8s.io/client-go v0.30.2/go.mod h1:JglKSWULm9xlJLx4KCkfLLQ7XwtlbflV6uFFSHTMgVs= k8s.io/klog/v2 v2.120.1 h1:QXU6cPEOIslTGvZaXvFWiP9VKyeet3sawzTOvdXb4Vw= k8s.io/klog/v2 v2.120.1/go.mod h1:3Jpz1GvMt720eyJH1ckRHK1EDfpxISzJ7I9OYgaDtPE= k8s.io/kube-openapi v0.0.0-20240228011516-70dd3763d340 h1:BZqlfIlq5YbRMFko6/PM7FjZpUb45WallggurYhKGag= diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 1aa83b26..45a4ef48 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -3,7 +3,6 @@ package main import ( "github.com/shuffle/shuffle-shared" - "bytes" "context" "encoding/json" @@ -22,18 +21,21 @@ import ( "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" + // This is for automatic removal of certain code :) "github.com/gorilla/mux" - "github.com/satori/go.uuid" + uuid "github.com/satori/go.uuid" //k8s deps corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + appsv1 "k8s.io/api/apps/v1" + "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/client-go/kubernetes" ) @@ -77,6 +79,7 @@ var startAction string //var allLogs map[string]string //var containerIds []string var downloadedImages []string + type ImageDownloadBody struct { Image string `json:"image"` } @@ -87,7 +90,7 @@ type ImageRequest struct { var finishedExecutions []string var imagesDistributed []string - +var imagedownloadTimeout = time.Second * 300 // Images to be autodeployed in the latest version of Shuffle. var autoDeploy = map[string]string{ @@ -175,7 +178,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 { @@ -183,19 +186,17 @@ 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 { @@ -214,9 +215,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 } } @@ -235,21 +236,20 @@ 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,7 +271,6 @@ 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) @@ -311,7 +310,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason } */ } else { - + } if len(reason) > 0 && len(nodeId) > 0 { @@ -369,7 +368,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason //Finished shutdown (after %d seconds). ", sleepDuration) // Allows everything to finish in subprocesses (apps) - if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" { + if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" && isKubernetes != "true" { time.Sleep(time.Duration(sleepDuration) * time.Second) os.Exit(3) } else { @@ -377,37 +376,46 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason } } -// 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 { - if isKubernetes == "true" { - if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 { - kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") - } else { - kubernetesNamespace = "default" +func int32Ptr(i int32) *int32 { return &i } + +// ** STARTREMOVE ***/ +func deployk8sApp(image string, identifier string, env []string) error { + if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 { + kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") + } else { + kubernetesNamespace = "default" + } + + envMap := make(map[string]string) + for _, envStr := range env { + parts := strings.SplitN(envStr, "=", 2) + if len(parts) == 2 { + envMap[parts[0]] = parts[1] } + } - envMap := make(map[string]string) - for _, envStr := range env { - parts := strings.SplitN(envStr, "=", 2) - if len(parts) == 2 { - envMap[parts[0]] = parts[1] - } - } + // add to env + // fmt.Sprintf("SHUFFLE_APP_EXPOSED_PORT=%d", deployport), + // fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", os.Getenv("SHUFFLE_SWARM_CONFIG")), + envMap["SHUFFLE_APP_EXPOSED_PORT"] = "80" + envMap["SHUFFLE_SWARM_CONFIG"] = os.Getenv("SHUFFLE_SWARM_CONFIG") + envMap["BASE_URL"] = "http://shuffle-workers:33333" - clientset, _, err := shuffle.GetKubernetesClient() - if err != nil { - log.Printf("[ERROR] Failed getting kubernetes: %s", err) - return err - } + clientset, _, err := shuffle.GetKubernetesClient() + if err != nil { + log.Printf("[ERROR] Failed getting kubernetes: %s", err) + return err + } - str := strings.ToLower(identifier) - strSplit := strings.Split(str, "_") - value := strSplit[0] - value = strings.ReplaceAll(value, "_", "-") + // str := strings.ToLower(identifier) + // strSplit := strings.Split(str, "_") + // value := strSplit[0] + // value = strings.ReplaceAll(value, "_", "-") + value := identifier - // Checking if app is generated or not - localRegistry := os.Getenv("REGISTRY_URL") - /* + // Checking if app is generated or not + localRegistry := os.Getenv("REGISTRY_URL") + /* appDetails := strings.Split(image, ":")[1] appDetailsSplit := strings.Split(appDetails, "_") appName := strings.Join(appDetailsSplit[:len(appDetailsSplit)-1], "_") @@ -425,63 +433,278 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } } } - */ + */ - if len(localRegistry) == 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")) > 0 { - localRegistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY") + if len(localRegistry) == 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")) > 0 { + localRegistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY") + } + + if len(localRegistry) > 0 && strings.Count(image, "/") <= 2 { + log.Printf("[DEBUG] Using REGISTRY_URL %s", localRegistry) + image = fmt.Sprintf("%s/%s", localRegistry, image) + } else { + if strings.Count(image, "/") <= 2 { + image = fmt.Sprintf("frikky/shuffle:%s", image) } + } - if len(localRegistry) > 0 && strings.Count(image, "/") <= 2 { - log.Printf("[DEBUG] Using REGISTRY_URL %s", localRegistry) - image = fmt.Sprintf("%s/%s", localRegistry, image) + log.Printf("[DEBUG] Got kubernetes with namespace %#v to run image '%s'", kubernetesNamespace, image) + + //fix naming convention + // podUuid := uuid.NewV4().String() + // podName := fmt.Sprintf("%s-%s", value, podUuid) + // replace identifier "_" with "-" + podName := strings.ReplaceAll(identifier, "_", "-") + + // pod := &corev1.Pod{ + // ObjectMeta: metav1.ObjectMeta{ + // Name: podName, + // Labels: map[string]string{ + // "app": podName, + // // "executionId": workflowExecution.ExecutionId, + // }, + // }, + // Spec: corev1.PodSpec{ + // RestartPolicy: "Never", // As a crash is not useful in this context + // // DNSPolicy: "Default", + // DNSPolicy: corev1.DNSClusterFirst, + // // NodeName: "worker1" + // Containers: []corev1.Container{ + // { + // Name: value, + // Image: image, + // Env: buildEnvVars(envMap), + + // // Pull if not available + // ImagePullPolicy: corev1.PullIfNotPresent, + // }, + // }, + // }, + // } + + // createdPod, err := clientset.CoreV1().Pods(kubernetesNamespace).Create(context.Background(), pod, metav1.CreateOptions{}) + // if err != nil { + // log.Printf("[ERROR] Failed creating pod: %v", err) + // // os.Exit(1) + // } else { + // log.Printf("[DEBUG] Created pod %#v in namespace %#v", createdPod.Name, kubernetesNamespace) + // } + + // service := &corev1.Service{ + // ObjectMeta: metav1.ObjectMeta{ + // Name: identifier, + // }, + // Spec: corev1.ServiceSpec{ + // Selector: map[string]string{ + // "app": podName, + // }, + // Ports: []corev1.ServicePort{ + // { + // Protocol: "TCP", + // Port: 80, + // TargetPort: intstr.FromInt(80), + // }, + // }, + // Type: corev1.ServiceTypeNodePort, + // }, + // } + + // _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + // if err != nil { + // log.Printf("[ERROR] Failed creating service: %v", err) + // return err + // } + + // use deployment instead of pod + // then expose a service similarly. + // number of replicas can be set to os.Getenv("SHUFFLE_SCALE_REPLICAS") + replicaNumberStr := os.Getenv("SHUFFLE_SCALE_REPLICAS") + replicaNumber := 1 + if len(replicaNumberStr) > 0 { + tmpInt, err := strconv.Atoi(replicaNumberStr) + if err != nil { + log.Printf("[ERROR] %s is not a valid number for replication", replicaNumberStr) } else { - if strings.Count(image, "/") <= 2 { - image = fmt.Sprintf("frikky/shuffle:%s", image) - } + replicaNumber = tmpInt + } + } - log.Printf("[DEBUG] Got kubernetes with namespace %#v to run image '%s'", kubernetesNamespace, image) + replicaNumberInt32 := int32(replicaNumber) - //fix naming convention - podUuid := uuid.NewV4().String() - podName := fmt.Sprintf("%s-%s", value, podUuid) - - pod := &corev1.Pod{ - ObjectMeta: metav1.ObjectMeta{ - Name: podName, - Labels: map[string]string{ - "app": "shuffle-app", - "executionId": workflowExecution.ExecutionId, + deployment := &appsv1.Deployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: podName, + }, + Spec: appsv1.DeploymentSpec{ + Replicas: int32Ptr(replicaNumberInt32), + Selector: &metav1.LabelSelector{ + MatchLabels: map[string]string{ + "app": podName, }, }, - Spec: corev1.PodSpec{ - RestartPolicy: "Never", // As a crash is not useful in this context - DNSPolicy: "Default", - // NodeName: "worker1" - Containers: []corev1.Container{ - { - Name: value, - Image: image, - Env: buildEnvVars(envMap), - - // Pull if not available - ImagePullPolicy: corev1.PullIfNotPresent, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: map[string]string{ + "app": podName, + }, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: value, + Image: image, + Env: buildEnvVars(envMap), + }, }, }, }, - } - - createdPod, err := clientset.CoreV1().Pods(kubernetesNamespace).Create(context.Background(), pod, metav1.CreateOptions{}) - if err != nil { - log.Printf("[ERROR] Failed creating pod: %v", err) - // os.Exit(1) - } else { - log.Printf("[DEBUG] Created pod %#v in namespace %#v", createdPod.Name, kubernetesNamespace) - } - - return nil + }, } + _, err = clientset.AppsV1().Deployments(kubernetesNamespace).Create(context.Background(), deployment, metav1.CreateOptions{}) + if err != nil { + log.Printf("[ERROR] Failed creating deployment: %v", err) + return err + } + + // kubectl expose deployment {podName} --type=NodePort --port=80 --target-port=80 + service := &corev1.Service{ + ObjectMeta: metav1.ObjectMeta{ + Name: podName, + }, + Spec: corev1.ServiceSpec{ + Selector: map[string]string{ + "app": podName, + }, + Ports: []corev1.ServicePort{ + { + Protocol: "TCP", + Port: 80, + TargetPort: intstr.FromInt(80), + }, + }, + Type: corev1.ServiceTypeNodePort, + }, + } + + _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + if err != nil { + log.Printf("[ERROR] Failed creating service: %v", err) + return err + } + + return nil +} + +//** ENDREMOVE ***/ + +// 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 { + // if isKubernetes == "true" { + // if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 { + // kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") + // } else { + // kubernetesNamespace = "default" + // } + + // envMap := make(map[string]string) + // for _, envStr := range env { + // parts := strings.SplitN(envStr, "=", 2) + // if len(parts) == 2 { + // envMap[parts[0]] = parts[1] + // } + // } + + // clientset, _, err := shuffle.GetKubernetesClient() + // if err != nil { + // log.Printf("[ERROR] Failed getting kubernetes: %s", err) + // return err + // } + + // str := strings.ToLower(identifier) + // strSplit := strings.Split(str, "_") + // value := strSplit[0] + // value = strings.ReplaceAll(value, "_", "-") + + // // Checking if app is generated or not + // localRegistry := os.Getenv("REGISTRY_URL") + // /* + // appDetails := strings.Split(image, ":")[1] + // appDetailsSplit := strings.Split(appDetails, "_") + // appName := strings.Join(appDetailsSplit[:len(appDetailsSplit)-1], "_") + // appVersion := appDetailsSplit[len(appDetailsSplit)-1] + // 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) + // if app.AppName == appName && app.AppVersion == appVersion { + // if app.Generated == true { + // log.Printf("[DEBUG] Generated app, setting local registry") + // image = fmt.Sprintf("%s/%s", localRegistry, image) + // break + // } else { + // log.Printf("[DEBUG] Not generated app, setting shuffle registry") + // } + // } + // } + // */ + + // if len(localRegistry) == 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY")) > 0 { + // localRegistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY") + // } + + // if len(localRegistry) > 0 && strings.Count(image, "/") <= 2 { + // log.Printf("[DEBUG] Using REGISTRY_URL %s", localRegistry) + // image = fmt.Sprintf("%s/%s", localRegistry, image) + // } else { + // if strings.Count(image, "/") <= 2 { + // image = fmt.Sprintf("frikky/shuffle:%s", image) + // } + // } + + // log.Printf("[DEBUG] Got kubernetes with namespace %#v to run image '%s'", kubernetesNamespace, image) + + // //fix naming convention + // podUuid := uuid.NewV4().String() + // podName := fmt.Sprintf("%s-%s", value, podUuid) + + // pod := &corev1.Pod{ + // ObjectMeta: metav1.ObjectMeta{ + // Name: podName, + // Labels: map[string]string{ + // "app": "shuffle-app", + // "executionId": workflowExecution.ExecutionId, + // }, + // }, + // Spec: corev1.PodSpec{ + // RestartPolicy: "Never", // As a crash is not useful in this context + // // DNSPolicy: "Default", + // DNSPolicy: corev1.DNSClusterFirst, + // // NodeName: "worker1" + // Containers: []corev1.Container{ + // { + // Name: value, + // Image: image, + // Env: buildEnvVars(envMap), + + // // Pull if not available + // ImagePullPolicy: corev1.PullIfNotPresent, + // }, + // }, + // }, + // } + + // createdPod, err := clientset.CoreV1().Pods(kubernetesNamespace).Create(context.Background(), pod, metav1.CreateOptions{}) + // if err != nil { + // log.Printf("[ERROR] Failed creating pod: %v", err) + // // os.Exit(1) + // } else { + // log.Printf("[DEBUG] Created pod %#v in namespace %#v", createdPod.Name, kubernetesNamespace) + // } + + // return nil + // } + // form basic hostConfig ctx := context.Background() @@ -500,7 +723,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 { @@ -542,7 +765,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 { @@ -583,7 +806,6 @@ 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) @@ -618,8 +840,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } func cleanupKubernetesExecution(clientset *kubernetes.Clientset, workflowExecution shuffle.WorkflowExecution, namespace string) error { - - workerName := fmt.Sprintf("worker-%s", workflowExecution.ExecutionId) + // workerName := fmt.Sprintf("worker-%s", workflowExecution.ExecutionId) labelSelector := fmt.Sprintf("app=shuffle-app,executionId=%s", workflowExecution.ExecutionId) podList, err := clientset.CoreV1().Pods(namespace).List(context.TODO(), metav1.ListOptions{ @@ -637,11 +858,11 @@ func cleanupKubernetesExecution(clientset *kubernetes.Clientset, workflowExecuti 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) - } - log.Printf("[DEBUG] %s in namespace %s deleted.", workerName, 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) + // } + // log.Printf("[DEBUG] %s in namespace %s deleted.", workerName, namespace) return nil } @@ -824,6 +1045,48 @@ func removeIndex(s []string, i int) []string { func getWorkerURLs() ([]string, error) { workerUrls := []string{} + if isKubernetes == "true" { + workerUrls = append(workerUrls, "http://shuffle-workers:33333") + // workerUrls = append(workerUrls, "http://192.168.29.16:33333") + + // get service "shuffle-workers" "Endpoints" + // serviceName := "shuffle-workers" + // clientset, _, err := shuffle.GetKubernetesClient() + // if err != nil { + // log.Println("[ERROR] Failed to get Kubernetes client:", err) + // return workerUrls, err + // } + + // services, err := clientset.CoreV1().Services("default").List(context.Background(), metav1.ListOptions{}) + // if err != nil { + // log.Println("[ERROR] Failed to list services:", err) + // return workerUrls, err + // } + + // for _, service := range services.Items { + // if service.Name == serviceName { + // endpoints, err := clientset.CoreV1().Endpoints("default").Get(context.Background(), serviceName, metav1.GetOptions{}) + // if err != nil { + // log.Println("[ERROR] Failed to get endpoints for service:", err) + // return workerUrls, err + // } + + // for _, subset := range endpoints.Subsets { + // for _, address := range subset.Addresses { + // for _, port := range subset.Ports { + // url := fmt.Sprintf("http://%s:%d", address.IP, port.Port) + // workerUrls = append(workerUrls, url) + // } + // } + // } + // } + // } + + log.Printf("[DEBUG] Worker URLs for k8s: %#v", workerUrls) + + return workerUrls, nil + } + // Create a new Docker client cli, err := dockerclient.NewEnvClient() if err != nil { @@ -863,7 +1126,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) { @@ -876,13 +1139,12 @@ func askOtherWorkersToDownloadImage(image string) { return } - if len(urls) < 2{ + if len(urls) < 2 { return } - httpClient := &http.Client{} - distributed := false + distributed := false for _, url := range urls { //log.Printf("[DEBUG] Trying to speak to: %s", url) imagesRequest := ImageRequest{ @@ -896,7 +1158,7 @@ func askOtherWorkersToDownloadImage(image string) { req, err := http.NewRequest( "POST", url, - bytes.NewBuffer(imageJSON), + bytes.NewBuffer(imageJSON), ) if err != nil { @@ -936,16 +1198,22 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { return } - startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId) - dockercli, err := dockerclient.NewEnvClient() - if err != nil { - log.Printf("[ERROR] Unable to create docker client (3): %s", err) - return + var dockercli *dockerclient.Client + var err error + + if isKubernetes != "true" { + dockercli, err = dockerclient.NewEnvClient() + if err != nil { + log.Printf("[ERROR] Unable to create docker client (3): %s", err) + return + } } - defer dockercli.Close() + if isKubernetes == "true" { + defer dockercli.Close() + } for _, action := range relevantActions { appname := action.AppName @@ -976,22 +1244,24 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { //executed = append(executed, action.ID) // FIXME - check whether it's running locally yet too - - stats, err := dockercli.ContainerInspect(context.Background(), identifier) - if err != nil || stats.ContainerJSONBase.State.Status != "running" { - // REMOVE - if err == nil { - log.Printf("[DEBUG][%s] Docker Container Status: %s, should kill: %s", workflowExecution.ExecutionId, stats.ContainerJSONBase.State.Status, identifier) - err = removeContainer(identifier) - if err != nil { - log.Printf("Error killing container: %s", err) + // take care of auto clean up later on for k8s + if isKubernetes != "true" { + stats, err := dockercli.ContainerInspect(context.Background(), identifier) + if err != nil || stats.ContainerJSONBase.State.Status != "running" { + // REMOVE + if err == nil { + log.Printf("[DEBUG][%s] Docker Container Status: %s, should kill: %s", workflowExecution.ExecutionId, stats.ContainerJSONBase.State.Status, identifier) + err = removeContainer(identifier) + if err != nil { + log.Printf("[ERROR] Error killing container: %s", err) + } + } else { + //log.Printf("WHAT TO DO HERE?: %s", err) } - } else { - //log.Printf("WHAT TO DO HERE?: %s", err) + } else if stats.ContainerJSONBase.State.Status == "running" { + //log.Printf(" + continue } - } else if stats.ContainerJSONBase.State.Status == "running" { - //log.Printf(" - continue } if len(action.Parameters) == 0 { @@ -1004,7 +1274,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) @@ -1090,10 +1360,9 @@ 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 @@ -1121,8 +1390,6 @@ 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" { @@ -1134,7 +1401,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { return } - err := downloadDockerImageBackend(&http.Client{Timeout: 60 * time.Second}, image) + err := downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) executed := false if err == nil { log.Printf("[DEBUG] Downloaded image %s from backend (CLEANUP)", image) @@ -1247,13 +1514,15 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { } log.Printf("[DEBUG][%s] Failed deploy. Downloading image %s: %s", workflowExecution.ExecutionId, image, err) - err := downloadDockerImageBackend(&http.Client{Timeout: 60 * time.Second}, image) + err := downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) + executed := false if err == nil { log.Printf("[DEBUG] Downloaded image %s from backend (CLEANUP)", image) //err = deployApp(dockercli, image, identifier, env, workflow, action) err = deployApp(dockercli, image, identifier, env, workflowExecution, action) if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") { + log.Printf("[ERROR] Err: %s", err) if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (40)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) @@ -1268,6 +1537,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { image = images[2] err = deployApp(dockercli, image, identifier, env, workflowExecution, action) if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") { + log.Printf("[ERROR] Err: %s", err) if strings.Contains(err.Error(), "exited prematurely") { log.Printf("[DEBUG] Shutting down (11)") shutdown(workflowExecution, action.ID, fmt.Sprintf("%s", err.Error()), true) @@ -1276,6 +1546,11 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { log.Printf("[WARNING] Failed deploying image THREE TIMES. Attempting to download %s as last resort from backend and dockerhub: %s", image, err) + if isKubernetes == "true" { + log.Printf("[ERROR] Image %s doesn't exist. Returning error for now") + return + } + reader, err := dockercli.ImagePull(context.Background(), image, pullOptions) if err != nil && !strings.Contains(err.Error(), "Conflict. The container name") { log.Printf("[ERROR] Failed getting %s. The couldn't be find locally, AND is missing.", image) @@ -1625,15 +1900,13 @@ 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) } - log.Printf("[INFO][%s] (2) Status: %s, Results: %d, actions: %d. Userinput: %#v", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra, hasUserinput) - return errors.New("Subflow status not found yet") + return errors.New("Subflow status not found yet") } func handleDefaultExecutionWrapper(ctx context.Context, workflowExecution shuffle.WorkflowExecution, streamResultUrl string, extra int) error { @@ -1898,7 +2171,6 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { if request.Body == nil { - log.Printf("[WARNING] (2) No body in request for workflowqueue") resp.WriteHeader(http.StatusBadRequest) return } @@ -1974,9 +2246,9 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { 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) + // 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) } @@ -2003,10 +2275,9 @@ 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) @@ -2055,38 +2326,38 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl //log.Printf(`[DEBUG][%s] Got result %s from %s. Execution status: %s. Save: %#v. Parent: %#v`, actionResult.ExecutionId, actionResult.Status, actionResult.Action.ID, workflowExecution.Status, dbSave, workflowExecution.ExecutionParent) //dbSave := false - + //if len(results) != len(workflowExecution.Results) { // log.Printf("[DEBUG][%s] There may have been an issue in transaction queue. Result lengths: %d vs %d. Should check which exists the base results, but not in entire execution, then append.", workflowExecution.ExecutionId, len(results), len(workflowExecution.Results)) //} - + // Validating that action results hasn't changed // Handled using cachhing, so actually pretty fast cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId) cache, err := shuffle.GetCache(ctx, cacheKey) if err == nil { //parsedValue := value.(*shuffle.WorkflowExecution) - + parsedValue := &shuffle.WorkflowExecution{} cacheData := []byte(cache.([]uint8)) err = json.Unmarshal(cacheData, &workflowExecution) if err != nil { 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 { } - + 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 + } */ } } @@ -2100,64 +2371,64 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed setting workflowexecution actionresult: %s"}`, err))) return } - + } else { log.Printf("[INFO][%s] Skipping setexec with status %s", workflowExecution.ExecutionId, workflowExecution.Status) - + // 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) } - + //if newExecutions && len(nextActions) > 0 { // log.Printf("[DEBUG][%s] New execution: %#v. NextActions: %#v", newExecutions, nextActions) // //handleExecutionResult(*workflowExecution) //} - + resp.WriteHeader(200) resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) } 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) return } - + if actionResult.ExecutionId == "TBD" { return } - + log.Printf("[DEBUG][%s] Sending FAILURE to self to stop the workflow execution. Action: %s (%s), app %s:%s", actionResult.ExecutionId, actionResult.Action.Label, actionResult.Action.ID, actionResult.Action.AppName, actionResult.Action.AppVersion) - + // Literally sending to same worker to run it as a new request streamUrl := fmt.Sprintf("http://localhost:33333/api/v1/streams") hostenv := os.Getenv("WORKER_HOSTNAME") if len(hostenv) > 0 { streamUrl = fmt.Sprintf("http://%s:33333/api/v1/streams", hostenv) } - + req, err := http.NewRequest( "POST", streamUrl, bytes.NewBuffer([]byte(data)), ) - + if err != nil { log.Printf("[ERROR][%s] Failed creating self request (1): %s", actionResult.ExecutionId, err) return } - + client := shuffle.GetExternalClient(streamUrl) newresp, err := client.Do(req) if err != nil { log.Printf("[ERROR][%s] Error running finishing request (2): %s", actionResult.ExecutionId, err) return } - + defer newresp.Body.Close() if newresp.Body != nil { body, err := ioutil.ReadAll(newresp.Body) @@ -2176,39 +2447,39 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { //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 - } + 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) + 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", streamUrl, bytes.NewBuffer([]byte(data)), ) - + if err != nil { log.Printf("[ERROR][%s] Failed creating finishing request: %s", workflowExecution.ExecutionId, err) log.Printf("[DEBUG][%s] Shutting down (22)", workflowExecution.ExecutionId) shutdown(workflowExecution, "", "", false) return } - + client := shuffle.GetExternalClient(streamUrl) newresp, err := client.Do(req) if err != nil { @@ -2217,7 +2488,7 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { shutdown(workflowExecution, "", "", false) return } - + defer newresp.Body.Close() if newresp.Body != nil { body, err := ioutil.ReadAll(newresp.Body) @@ -2228,11 +2499,11 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { log.Printf("[DEBUG][%s] NEWRESP (from backend): %s", workflowExecution.ExecutionId, string(body)) } } - } - - func validateFinished(workflowExecution shuffle.WorkflowExecution) bool { +} + +func validateFinished(workflowExecution shuffle.WorkflowExecution) bool { ctx := context.Background() - + newexec, err := shuffle.GetWorkflowExecution(ctx, workflowExecution.ExecutionId) if err != nil { log.Printf("[ERROR][%s] Failed getting workflow execution: %s", workflowExecution.ExecutionId, err) @@ -2240,15 +2511,15 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { } else { workflowExecution = *newexec } - + //startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId) 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", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Workflow.Actions), extra, len(workflowExecution.Results), workflowExecution.ExecutionParent) - - 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" || 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" { @@ -2257,17 +2528,17 @@ 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) shutdownData, err := json.Marshal(workflowExecution) if err != nil { log.Printf("[ERROR][%s] Shutting down (32): Failed to unmarshal data for backend: %s", workflowExecution.ExecutionId, err) shutdown(workflowExecution, "", "", true) } - + cacheKey := fmt.Sprintf("workflowexecution_%s", workflowExecution.ExecutionId) if len(workflowExecution.Authorization) > 0 { err = shuffle.SetCache(ctx, cacheKey, shutdownData, 31) @@ -2275,12 +2546,12 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) { log.Printf("[ERROR][%s] Failed adding to cache during ValidateFinished", workflowExecution) } } - + shuffle.RunCacheCleanup(ctx, workflowExecution) sendResult(workflowExecution, shutdownData) return true } - + return false } @@ -2293,7 +2564,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return } - + var actionResult shuffle.ActionResult err = json.Unmarshal(body, &actionResult) if err != nil { @@ -2302,14 +2573,14 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { //resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) //return } - + if len(actionResult.ExecutionId) == 0 { log.Printf("[WARNING] No workflow execution id in action result (2). Data: %s", string(body)) resp.WriteHeader(400) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "No workflow execution id in action result"}`))) return } - + ctx := context.Background() workflowExecution, err := shuffle.GetWorkflowExecution(ctx, actionResult.ExecutionId) if err != nil { @@ -2318,7 +2589,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key or execution_id might not exist."}`))) return } - + // Authorization is done here if workflowExecution.Authorization != actionResult.Authorization { log.Printf("[ERROR] Bad authorization key when getting stream results from cache %s.", actionResult.ExecutionId) @@ -2326,14 +2597,14 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Bad authorization key or execution_id might not exist."}`))) return } - + newjson, err := json.Marshal(workflowExecution) if err != nil { resp.WriteHeader(500) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed unpacking workflow execution"}`))) return } - + resp.WriteHeader(200) resp.Write(newjson) @@ -2342,12 +2613,12 @@ 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 "" } - + for _, address := range addrs { // check the address type and if it is not a loopback the display it if ipnet, ok := address.(*net.IPNet); ok && !ipnet.IP.IsLoopback() { @@ -2356,7 +2627,7 @@ func getLocalIP() string { } } } - + return "" } @@ -2367,17 +2638,21 @@ func getAvailablePort() (net.Listener, error) { //return ":5001" return nil, err } - + //defer listener.Close() - + return listener, nil //return fmt.Sprintf(":%d", port) } func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { hostname = getLocalIP() - os.Setenv("WORKER_HOSTNAME", hostname) - + if isKubernetes == "true" { + os.Setenv("WORKER_HOSTNAME", "shuffle-workers") + } else { + os.Setenv("WORKER_HOSTNAME", hostname) + } + // FIXME: This MAY not work because of speed between first // container being launched and port being assigned to webserver listener, err := getAvailablePort() @@ -2385,17 +2660,17 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { log.Printf("[ERROR] Failed to create init listener: %s", err) return 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)) - + 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) return listener } @@ -2406,25 +2681,25 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { 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 - 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) - + req, err := http.NewRequest( "POST", dockerImgUrl, bytes.NewBuffer([]byte(data)), ) - + authorization := os.Getenv("AUTHORIZATION") if len(authorization) > 0 { req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization)) @@ -2432,28 +2707,28 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { log.Printf("[WARNING] No auth found - running backend download without it.") //return } - + newresp, err := topClient.Do(req) if err != nil { log.Printf("[ERROR] Failed download request for %s: %s", imageName, err) return err } - + defer newresp.Body.Close() if newresp.StatusCode != 200 { log.Printf("[ERROR] Docker download for image %s (backend) StatusCode (1): %d", imageName, newresp.StatusCode) return errors.New(fmt.Sprintf("Failed to get image - status code %d", newresp.StatusCode)) } - + newImageName := strings.Replace(imageName, "/", "_", -1) newFileName := newImageName + ".tar" - + tar, err := os.Create(newFileName) if err != nil { log.Printf("[WARNING] Failed creating file: %s", err) return err } - + defer tar.Close() _, err = io.Copy(tar, newresp.Body) if err != nil { @@ -2461,60 +2736,60 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { return err } tar.Seek(0, 0) - + dockercli, err := dockerclient.NewEnvClient() if err != nil { log.Printf("[ERROR] Unable to create docker client (3): %s", err) return err } - + defer dockercli.Close() - + imageLoadResponse, err := dockercli.ImageLoad(context.Background(), tar, true) if err != nil { log.Printf("[ERROR] Error loading images: %s", err) return err } - + defer imageLoadResponse.Body.Close() body, err := ioutil.ReadAll(imageLoadResponse.Body) if err != nil { log.Printf("[ERROR] Error reading: %s", err) return err } - + if strings.Contains(string(body), "no such file") { return errors.New(string(body)) } - + baseTag := strings.Split(imageName, ":") if len(baseTag) > 1 { tag := baseTag[1] log.Printf("[DEBUG] Creating tag copies of downloaded containers from tag %s", tag) - + // Remapping ctx := context.Background() dockercli.ImageTag(ctx, imageName, fmt.Sprintf("frikky/shuffle:%s", tag)) dockercli.ImageTag(ctx, imageName, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", tag)) - + downloadedImages = append(downloadedImages, fmt.Sprintf("frikky/shuffle:%s", tag)) downloadedImages = append(downloadedImages, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", tag)) - + } - + os.Remove(newFileName) - + log.Printf("[INFO] Successfully loaded image %s: %s", imageName, string(body)) return nil - } - - func findActiveSwarmNodes(dockercli *dockerclient.Client) (int64, error) { +} + +func findActiveSwarmNodes(dockercli *dockerclient.Client) (int64, error) { ctx := context.Background() nodes, err := dockercli.NodeList(ctx, types.NodeListOptions{}) if err != nil { return 1, err } - + nodeCount := int64(0) for _, node := range nodes { //log.Printf("ID: %s - %#v", node.ID, node.Status.State) @@ -2522,7 +2797,7 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { nodeCount += 1 } } - + // Check for SHUFFLE_MAX_NODES maxNodesString := os.Getenv("SHUFFLE_MAX_SWARM_NODES") // Make it into a number and check if it's lower than nodeCount @@ -2531,14 +2806,14 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error { if err != nil { return nodeCount, err } - + if nodeCount > maxNodes { nodeCount = maxNodes } } - + return nodeCount, nil - + /* containers, err := dockercli.ContainerList(ctx, types.ContainerListOptions{ All: true, @@ -2550,156 +2825,160 @@ 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 { -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) + 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 -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) + // Specific for subflow to ensure worker matches the backend correctly - //parsedRequest.Url + 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) + } } - //log.Printf("[DEBUG][%s] Should add a baseurl for the app to get back to: %s", workflowExecution.ExecutionId, parsedRequest.Url) -} + 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) -// Swapping because this was confusing during dev -// No real reason, just variable names -tmp := parsedRequest.Url -parsedRequest.Url = parsedRequest.BaseUrl -parsedRequest.BaseUrl = tmp + //parsedRequest.Url + } -// 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) + //log.Printf("[DEBUG][%s] Should add a baseurl for the app to get back to: %s", workflowExecution.ExecutionId, parsedRequest.Url) + } - if parsedRequest.Action.AppName == "shuffle-subflow" || parsedRequest.Action.AppName == "shuffle-subflow-v2" || parsedRequest.Action.AppName == "User Input" { + // 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.Url = parsedRequest.BaseUrl + //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 -} + // 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) -} - -client := shuffle.GetExternalClient(streamUrl) -customTimeout := os.Getenv("SHUFFLE_APP_REQUEST_TIMEOUT") -if len(customTimeout) > 0 { - // convert to int - timeoutInt, err := strconv.Atoi(customTimeout) + data, err := json.Marshal(parsedRequest) if err != nil { - log.Printf("[ERROR] Failed converting SHUFFLE_APP_REQUEST_TIMEOUT to int: %s", err) - } else { - log.Printf("[DEBUG] Setting client timeout to %d seconds for app request", timeoutInt) - client.Timeout = time.Duration(timeoutInt) * time.Second + log.Printf("[ERROR] Failed marshalling worker request: %s", err) + return err } -} -newresp, err := client.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") { + if isKubernetes == "true" { + appName = strings.Replace(appName, "_", "-", -1) + } + + 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 } - 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) + 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 { - // escape quotes and newlines - newerr = strings.ReplaceAll(strings.ReplaceAll(newerr, "\"", "\\\""), "\n", "\\n") + //log.Printf("[DEBUG][%s] Adding %s to cache (%#v)", workflowExecution.ExecutionId, newExecId, action.Name) } - if strings.Contains(fmt.Sprintf("%s", err), "no such host") { - log.Printf("[DEBUG] SHOULD be Removing references to location for app %s as to be rediscovered", action.AppName) - - //for k, v := range portMappings { - // if strings.Contains(strings.ToLower(strings.ReplaceAll(action.AppName, " ", "_"))) { - // } - //} - - //var portMappings map[string]int + client := shuffle.GetExternalClient(streamUrl) + customTimeout := os.Getenv("SHUFFLE_APP_REQUEST_TIMEOUT") + if len(customTimeout) > 0 { + // convert to int + timeoutInt, err := strconv.Atoi(customTimeout) + if err != nil { + log.Printf("[ERROR] Failed converting SHUFFLE_APP_REQUEST_TIMEOUT to int: %s", err) + } else { + log.Printf("[DEBUG] Setting client timeout to %d seconds for app request", timeoutInt) + client.Timeout = time.Duration(timeoutInt) * time.Second + } } - 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. Try the action again, 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()), + newresp, err := client.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") + } + + if strings.Contains(fmt.Sprintf("%s", err), "no such host") { + log.Printf("[DEBUG] SHOULD be Removing references to location for app %s as to be rediscovered", action.AppName) + + //for k, v := range portMappings { + // if strings.Contains(strings.ToLower(strings.ReplaceAll(action.AppName, " ", "_"))) { + // } + //} + + //var portMappings map[string]int + } + + 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. Try the action again, 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", } @@ -2724,13 +3003,17 @@ if err != nil { // 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 - } + var cli *dockerclient.Client + var err error - defer cli.Close() + if isKubernetes != "true" { + cli, err := dockerclient.NewEnvClient() + if err != nil { + log.Printf("[ERROR] Unable to create docker client (3): %s", err) + return + } + defer cli.Close() + } for key, value := range autoDeploy { newNameSplit := strings.Split(key, ":") @@ -2810,7 +3093,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) @@ -2850,7 +3133,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe } // 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{}) + childNodes := shuffle.FindChildNodes(workflowExecution.Workflow, 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 { @@ -2893,7 +3176,6 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe // Set environment variable - //log.Printf("Before wait") //wg := sync.WaitGroup{} //wg.Add(1) @@ -2944,7 +3226,11 @@ func main() { 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) + if isKubernetes != "true" { + log.Printf("[DEBUG] Ran init for worker to set up cache system. Docker version: %s", dockerApiVersion) + } else { + log.Printf("[DEBUG] Ran init for worker to set up cache system on Kubernetes") + } } log.Printf("[INFO] Setting up worker environment") @@ -3270,7 +3556,6 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) { return } - for _, img := range images { for _, tag := range img.RepoTags { splitTag := strings.Split(tag, ":") @@ -3283,7 +3568,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"}`))) @@ -3293,7 +3578,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) { } log.Printf("[INFO] Downloading image %s", image.Image) - downloadDockerImageBackend(&http.Client{Timeout: 60 * time.Second}, image.Image) + downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image.Image) // return success resp.WriteHeader(200) @@ -3326,6 +3611,7 @@ func runWebserver(listener net.Listener) { //log.Fatal(http.Serve(listener, nil)) + log.Printf("[DEBUG] NEW webserver setup") http.Handle("/", r)