From 26ee1d4da57c95595726b9b3c3b07135bbeff25f Mon Sep 17 00:00:00 2001 From: Frikky Date: Tue, 11 Jun 2024 16:05:15 +0200 Subject: [PATCH] Minor sdk fix --- backend/app_sdk/app_base.py | 22 +++++++++- backend/go-app/main.go | 3 +- backend/go-app/walkoff.go | 6 +++ functions/onprem/orborus/go.mod | 2 +- functions/onprem/orborus/go.sum | 2 + functions/onprem/orborus/orborus.go | 63 +++++++++++++++++++++++------ 6 files changed, 82 insertions(+), 16 deletions(-) diff --git a/backend/app_sdk/app_base.py b/backend/app_sdk/app_base.py index 6ef571aa..a7e0d4f8 100755 --- a/backend/app_sdk/app_base.py +++ b/backend/app_sdk/app_base.py @@ -652,7 +652,14 @@ class AppBase: sleeptime = float(random.randint(0, 10) / 10) try: - ret = requests.post(url, headers=headers, json=action_result, timeout=10, verify=False, proxies=self.proxy_config) + ret = requests.post( + url, + headers=headers, + json=action_result, + timeout=10, + verify=False, + proxies=self.proxy_config, + ) #self.logger.info(f"""[DEBUG] Successful result request: Status= {ret.status_code} (break on 200/201) & Action status: {action_result["status"]}. Response= {ret.text}""") if ret.status_code == 200 or ret.status_code == 201: @@ -660,7 +667,18 @@ class AppBase: break else: # FIXME: Add a checker for 403, and Proxy logs failing - self.logger.info(f"[ERROR] Bad resp ({ret.status_code}) in send_result for url '{url}'") + headerauth = "" + if "Authorization" in headers: + headerauth = headers["Authorization"] + + try: + + self.logger.info(f"[ERROR] Bad resp ({ret.status_code}) in send_result for url '{url}'. Execution ID: %d, Authorization: %d, Header Auth: %d" % (len(action_result["execution_id"]), len(action_result["authorization"]), len(headerauth))) + + except Exception as e: + self.logger.info(f"[ERROR] Bad resp ({ret.status_code}) in send_result for url '{url}' (no detail)") + pass + time.sleep(sleeptime) diff --git a/backend/go-app/main.go b/backend/go-app/main.go index 98718c19..d51535fa 100755 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -5017,6 +5017,7 @@ func initHandlers() { r.HandleFunc("/api/v1/workflows/{key}/schedule/{schedule}", stopSchedule).Methods("DELETE", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}/stream", shuffle.HandleStreamWorkflow).Methods("GET", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}/stream", shuffle.HandleStreamWorkflowUpdate).Methods("POST", "OPTIONS") + r.HandleFunc("/api/v1/workflows/{key}/duplicate", shuffle.DuplicateWorkflow).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}", deleteWorkflow).Methods("DELETE", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}", shuffle.SaveWorkflow).Methods("PUT", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}", shuffle.GetSpecificWorkflow).Methods("GET", "OPTIONS") @@ -5029,6 +5030,7 @@ func initHandlers() { r.HandleFunc("/api/v1/recommendations/get_actions", shuffle.HandleActionRecommendation).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/recommendations/modify", shuffle.HandleRecommendationAction).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/workflows/{key}/revisions", shuffle.GetWorkflowRevisions).Methods("GET", "OPTIONS") + r.HandleFunc("/api/v1/workflows/{key}/child_workflows", shuffle.GetChildWorkflows).Methods("GET", "OPTIONS") // Triggers r.HandleFunc("/api/v1/hooks/new", shuffle.HandleNewHook).Methods("POST", "OPTIONS") @@ -5111,7 +5113,6 @@ func initHandlers() { // Important for email, IDS etc. Create this by: // PS: For cloud, this has to use cloud storage. // https://developer.box.com/reference/get-files-id-content/ - // 1. Creating the "get file" option. Make it possible to run this in the frontend. r.HandleFunc("/api/v1/files/download_remote", shuffle.HandleDownloadRemoteFiles).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/files/namespaces/{namespace}", shuffle.HandleGetFileNamespace).Methods("GET", "OPTIONS") r.HandleFunc("/api/v1/files/{fileId}/content", shuffle.HandleGetFileContent).Methods("GET", "OPTIONS") diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 10df25b2..b6b3258c 100755 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -921,6 +921,12 @@ func deleteWorkflow(resp http.ResponseWriter, request *http.Request) { return } + if len(workflow.ParentWorkflowId) > 0 { + resp.WriteHeader(403) + resp.Write([]byte(`{"success": false, "reason": "Can't delete a workflow distributed from your parent org"}`)) + return + } + if user.Id != workflow.Owner || len(user.Id) == 0 { if workflow.OrgId == user.ActiveOrg.Id && user.Role == "admin" { log.Printf("[INFO] User %s is deleting workflow %s as admin. Owner: %s", user.Username, workflow.ID, workflow.Owner) diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod index 9101801f..8faab828 100644 --- a/functions/onprem/orborus/go.mod +++ b/functions/onprem/orborus/go.mod @@ -4,7 +4,7 @@ go 1.22.0 toolchain go1.22.2 -//replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared +replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared require ( github.com/docker/docker v26.1.0+incompatible diff --git a/functions/onprem/orborus/go.sum b/functions/onprem/orborus/go.sum index 923438f9..1c27a7a1 100644 --- a/functions/onprem/orborus/go.sum +++ b/functions/onprem/orborus/go.sum @@ -393,6 +393,8 @@ github.com/shuffle/shuffle-shared v0.6.18 h1:mKc3vGuCz9ubdqMwaLocSbEUZyso620CY73 github.com/shuffle/shuffle-shared v0.6.18/go.mod h1:00QOcSPlUWMXzJj1D7pjcV9h6nVRWfSWqTM10+fhTd0= 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.37 h1:IB8tJqubJmJwwpLYbXNMVWck5jgVO9SKcKo/BBu+wTc= +github.com/shuffle/shuffle-shared v0.6.37/go.mod h1:rWkh1eWdIx7OqQzJ1+JzF3Hck1X/Ty1WkUtjLrp+CU4= 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/orborus/orborus.go b/functions/onprem/orborus/orborus.go index f9e28d44..eddcfd71 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -683,6 +683,7 @@ func deployWorker(image string, identifier string, env []string, executionReques return err } + //env = append(env, fmt.Sprintf("KUBERNETES_CONFIG=%s", config.String())) // FIXME: When a service account is used, the account is also mounted in the pod @@ -697,6 +698,7 @@ func deployWorker(image string, identifier string, env []string, executionReques // use k8s downward API to find it if we are in a pod } + // Check if namespace exist as variable. If so, make it if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 && !namespacemade { kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") @@ -720,6 +722,18 @@ func deployWorker(image string, identifier string, env []string, executionReques } } + if len(kubernetesNamespace) == 0 { + foundNamespace, err := shuffle.GetKubernetesNamespace() + if err != nil { + //log.Printf("[ERROR] Failed getting Kubernetes namespace: %s", err) + } + + if len(foundNamespace) > 0 { + kubernetesNamespace = foundNamespace + os.Setenv("KUBERNETES_NAMESPACE", kubernetesNamespace) + } + } + if len(kubernetesNamespace) == 0 { kubernetesNamespace = "default" } @@ -739,12 +753,40 @@ func deployWorker(image string, identifier string, env []string, executionReques } } + containerLabels := map[string]string{ + "container": "shuffle-worker", + } + + containerAttachment := corev1.Container{ + Name: identifier, + Image: kubernetesImage, + Env: buildEnvVars(envMap), + + //ImagePullPolicy: "Never", + ImagePullPolicy: corev1.PullIfNotPresent, + } + + podname := shuffle.GetPodName() + + ctx := context.Background() + + if len(podname) > 0 { + currentPodStatus, err := shuffle.GetCurrentPodNetworkConfig(ctx, clientset, kubernetesNamespace, podname) + if err != nil { + log.Printf("[ERROR] Failed getting current pod network: %s", err) + } else { + log.Printf("[DEBUG] Current pod found!") + // currentPodStatus = k8s.io/api/core/v1.PodStatus + } + } + + // While testing: // kubectl delete pods --all --all-namespaces; kubectl delete services --all --all-namespaces pod := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ Name: identifier, - Labels: map[string]string{"app": "shuffle-worker"}, + Labels: containerLabels, }, Spec: corev1.PodSpec{ RestartPolicy: "Never", @@ -753,23 +795,20 @@ func deployWorker(image string, identifier string, env []string, executionReques // "node": "master", // }, Containers: []corev1.Container{ - { - Name: identifier, - Image: kubernetesImage, - Env: buildEnvVars(envMap), - - //ImagePullPolicy: "Never", - ImagePullPolicy: corev1.PullIfNotPresent, - //ImagePullPolicy: "Always", - }, + containerAttachment, }, }, } // Check if running on ARM or x86 to download the correct image - // Add environment variables - // pod.Spec.Containers[0].Env = buildEnvVars(envMap) + // Get current pod's network so we can make the pod in it + + networks, err := clientset.CoreV1().Pods(kubernetesNamespace).List(context.Background(), metav1.ListOptions{}) + if err != nil { + log.Printf("[ERROR] Failed listing pods: %s", err) + } + createdPod, err := clientset.CoreV1().Pods(kubernetesNamespace).Create(context.Background(), pod, metav1.CreateOptions{}) if err != nil {