From 48505261983338d2845134f553d96fbef1fe7939 Mon Sep 17 00:00:00 2001 From: "trusihin.andrey" Date: Tue, 26 Aug 2025 15:09:39 +0300 Subject: [PATCH 1/4] Worker. Add option to set k8s resources. --- functions/onprem/orborus/orborus.go | 75 +++++++++++++++++++++++------ 1 file changed, 59 insertions(+), 16 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 26f16017..d8870c59 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -51,6 +51,7 @@ import ( appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" rbacv1 "k8s.io/api/rbac/v1" + "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" ) @@ -734,7 +735,7 @@ func deployServiceWorkers(image string) { var updatedNetworks []swarm.NetworkAttachmentConfig for _, net := range serviceSpec.Networks { if net.Target != "shuffle_shuffle" { - updatedNetworks = append(updatedNetworks, net) + updatedNetworks = append(updatedNetworks, net) } } serviceSpec.Networks = updatedNetworks @@ -786,6 +787,48 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { return envVars } +func buildResourcesFromEnv() corev1.ResourceRequirements { + reqs := corev1.ResourceList{} + lims := corev1.ResourceList{} + + type item struct { + env string + rn corev1.ResourceName + to *corev1.ResourceList + } + + items := []item{ + // kubernetes requests + {env: "KUBERNETES_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, + {env: "KUBERNETES_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, + {env: "KUBERNETES_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, + // kubernetes limits + {env: "KUBERNETES_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, + {env: "KUBERNETES_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, + {env: "KUBERNETES_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, + } + + for _, it := range items { + if v := strings.TrimSpace(os.Getenv(it.env)); v != "" { + if q, err := resource.ParseQuantity(v); err == nil { + (*it.to)[it.rn] = q + } else { + log.Printf("[WARN] Cannot parse %s=%q as resource quantity: %v", it.env, v, err) + } + } + } + + rr := corev1.ResourceRequirements{} + if len(reqs) > 0 { + rr.Requests = reqs + } + if len(lims) > 0 { + rr.Limits = lims + } + + return rr +} + func handleBackendImageDownload(ctx context.Context, images string) error { // Replicate images with lowercase, as the name may be wrong @@ -813,20 +856,20 @@ func handleBackendImageDownload(ctx context.Context, images string) error { newImages = append(newImages, curimage) // Force remove the current image to avoid cached layers - // if swarmConfig == "run" || swarmConfig == "swarm" { - // _, err := dockercli.ImageRemove(ctx, curimage, image.RemoveOptions{ - // Force: true, - // PruneChildren: true, - // }) + // if swarmConfig == "run" || swarmConfig == "swarm" { + // _, err := dockercli.ImageRemove(ctx, curimage, image.RemoveOptions{ + // Force: true, + // PruneChildren: true, + // }) - // if err != nil { - // log.Printf("[ERROR] Failed removing image for re-download: %s", err) - // } else { - // log.Printf("[DEBUG] Removed image: %s", curimage) - // } - // } else { - // //log.Printf("[DEBUG] Skipping image removal for %s as swarmConfig is not set to run or swarm. Value: %#v", curimage, swarmConfig) - // } + // if err != nil { + // log.Printf("[ERROR] Failed removing image for re-download: %s", err) + // } else { + // log.Printf("[DEBUG] Removed image: %s", curimage) + // } + // } else { + // //log.Printf("[DEBUG] Skipping image removal for %s as swarmConfig is not set to run or swarm. Value: %#v", curimage, swarmConfig) + // } err := shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, curimage) if err != nil { @@ -887,7 +930,7 @@ func handleBackendImageDownload(ctx context.Context, images string) error { log.Printf("[ERROR] Failed updating service %s with the new image %s: %s. Resp: %#v", service.Spec.Annotations.Name, image, err, resp) } else { log.Printf("[DEBUG] Updated service %s with the new image %s. Resp: %#v", service.Spec.Annotations.Name, image, resp) - + found = true if !strings.Contains(fmt.Sprintf("%s", resp), "error") { @@ -1173,6 +1216,7 @@ func deployK8sWorker(image string, identifier string, env []string) error { Image: kubernetesImage, Env: buildEnvVars(envMap), SecurityContext: containerSecurityContext, + Resources: buildResourcesFromEnv(), //ImagePullPolicy: "Never", ImagePullPolicy: corev1.PullIfNotPresent, @@ -2214,7 +2258,6 @@ func main() { log.Printf("[INFO] Waiting for executions at %s with Environment %#v", fullUrl, environment) - hasStarted := false for { if req.Method == "POST" { From f26f8728356173370de71465c22907e4b2a0fdbb Mon Sep 17 00:00:00 2001 From: "trusihin.andrey" Date: Tue, 26 Aug 2025 16:24:01 +0300 Subject: [PATCH 2/4] Rename k8s resources envs. Add the same options to shuffle app --- functions/onprem/orborus/orborus.go | 40 +++++++++++++++++---- functions/onprem/worker/worker.go | 55 ++++++++++++++++++++++++++--- 2 files changed, 84 insertions(+), 11 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index d8870c59..3103405a 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -799,13 +799,13 @@ func buildResourcesFromEnv() corev1.ResourceRequirements { items := []item{ // kubernetes requests - {env: "KUBERNETES_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, - {env: "KUBERNETES_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, - {env: "KUBERNETES_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, + {env: "SHUFFLE_WORKER_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, + {env: "SHUFFLE_WORKER_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, + {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, // kubernetes limits - {env: "KUBERNETES_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, - {env: "KUBERNETES_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, - {env: "KUBERNETES_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, + {env: "SHUFFLE_WORKER_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, + {env: "SHUFFLE_WORKER_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, + {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, } for _, it := range items { @@ -1072,6 +1072,34 @@ func deployK8sWorker(image string, identifier string, env []string) error { env = append(env, fmt.Sprintf("IS_KUBERNETES=true")) env = append(env, fmt.Sprintf("KUBERNETES_NAMESPACE=%s", os.Getenv("KUBERNETES_NAMESPACE"))) + // worker resource env + for _, k := range []string{ + "SHUFFLE_WORKER_CPU_REQUEST", + "SHUFFLE_WORKER_MEMORY_REQUEST", + "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", + "SHUFFLE_WORKER_CPU_LIMIT", + "SHUFFLE_WORKER_MEMORY_LIMIT", + "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", + } { + if v := os.Getenv(k); v != "" { + env = append(env, fmt.Sprintf("%s=%s", k, v)) + } + } + + // app resource env + for _, k := range []string{ + "SHUFFLE_APP_CPU_REQUEST", + "SHUFFLE_APP_MEMORY_REQUEST", + "SHUFFLE_APP_EPHEMERAL_STORAGE_REQUEST", + "SHUFFLE_APP_CPU_LIMIT", + "SHUFFLE_APP_MEMORY_LIMIT", + "SHUFFLE_APP_EPHEMERAL_STORAGE_LIMIT", + } { + if v := os.Getenv(k); v != "" { + env = append(env, fmt.Sprintf("%s=%s", k, v)) + } + } + if len(os.Getenv("KUBERNETES_SERVICE_HOST")) > 0 { env = append(env, fmt.Sprintf("KUBERNETES_SERVICE_HOST=%s", os.Getenv("KUBERNETES_SERVICE_HOST"))) } diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 2fbf51e3..133545f3 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -42,6 +42,7 @@ import ( //k8s deps appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/resource" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" "k8s.io/client-go/kubernetes" @@ -639,6 +640,7 @@ func deployk8sApp(image string, identifier string, env []string) error { }, }, SecurityContext: containerSecurityContext, + Resources: buildResourcesFromEnv(), }, }, DNSPolicy: corev1.DNSClusterFirst, @@ -2375,6 +2377,49 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { } return envVars } + +func buildResourcesFromEnv() corev1.ResourceRequirements { + reqs := corev1.ResourceList{} + lims := corev1.ResourceList{} + + type item struct { + env string + rn corev1.ResourceName + to *corev1.ResourceList + } + + items := []item{ + // kubernetes requests + {env: "SHUFFLE_APP_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, + {env: "SHUFFLE_APP_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, + {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, + // kubernetes limits + {env: "SHUFFLE_APP_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, + {env: "SHUFFLE_APP_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, + {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, + } + + for _, it := range items { + if v := strings.TrimSpace(os.Getenv(it.env)); v != "" { + if q, err := resource.ParseQuantity(v); err == nil { + (*it.to)[it.rn] = q + } else { + log.Printf("[WARN] Cannot parse %s=%q as resource quantity: %v", it.env, v, err) + } + } + } + + rr := corev1.ResourceRequirements{} + if len(reqs) > 0 { + rr.Requests = reqs + } + if len(lims) > 0 { + rr.Limits = lims + } + + return rr +} + func getWorkerBackendExecution(auth string, executionId string) (*shuffle.WorkflowExecution, error) { backendUrl := os.Getenv("BASE_URL") if len(backendUrl) == 0 { @@ -2385,9 +2430,9 @@ func getWorkerBackendExecution(auth string, executionId string) (*shuffle.Workfl streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", backendUrl) topClient := shuffle.GetExternalClient(backendUrl) - requestData := shuffle.ActionResult { + requestData := shuffle.ActionResult{ Authorization: auth, - ExecutionId: executionId, + ExecutionId: executionId, } data, err := json.Marshal(requestData) @@ -2647,7 +2692,7 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl } if setExecution || workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || workflowExecution.Status == "FAILURE" { - if debug { + if debug { log.Printf("[DEBUG][%s] Running setexec with status %s and %d/%d results", workflowExecution.ExecutionId, workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)) } @@ -2665,7 +2710,7 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl if os.Getenv("SHUFFLE_SWARM_CONFIG") == "run" || os.Getenv("SHUFFLE_SWARM_CONFIG") == "swarm" { finished := shuffle.ValidateFinished(ctx, -1, *workflowExecution) if !finished { - if debug { + if debug { log.Printf("[DEBUG][%s] Handling next node since it's not finished!", workflowExecution.ExecutionId) } @@ -3630,7 +3675,7 @@ func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, log.Printf("[ERROR] Failed reading app request body body: %s", err) return err } else { - if debug { + if debug { log.Printf("[DEBUG][%s] NEWRESP (from app): %s", workflowExecution.ExecutionId, string(body)) } } From dd89bbdd33f70625a0314a549c4169927d8b187a Mon Sep 17 00:00:00 2001 From: "trusihin.andrey" Date: Tue, 26 Aug 2025 18:10:12 +0300 Subject: [PATCH 3/4] Remove SHUFFLE_WORKER_* env from orborus.go --- functions/onprem/orborus/orborus.go | 14 -------------- 1 file changed, 14 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 3103405a..6d0d73ea 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -1072,20 +1072,6 @@ func deployK8sWorker(image string, identifier string, env []string) error { env = append(env, fmt.Sprintf("IS_KUBERNETES=true")) env = append(env, fmt.Sprintf("KUBERNETES_NAMESPACE=%s", os.Getenv("KUBERNETES_NAMESPACE"))) - // worker resource env - for _, k := range []string{ - "SHUFFLE_WORKER_CPU_REQUEST", - "SHUFFLE_WORKER_MEMORY_REQUEST", - "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", - "SHUFFLE_WORKER_CPU_LIMIT", - "SHUFFLE_WORKER_MEMORY_LIMIT", - "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", - } { - if v := os.Getenv(k); v != "" { - env = append(env, fmt.Sprintf("%s=%s", k, v)) - } - } - // app resource env for _, k := range []string{ "SHUFFLE_APP_CPU_REQUEST", From b6e90050371ef06b7c8ff95d93ef02f76a737a78 Mon Sep 17 00:00:00 2001 From: "trusihin.andrey" Date: Wed, 27 Aug 2025 09:26:47 +0300 Subject: [PATCH 4/4] Fix variables name --- functions/onprem/orborus/orborus.go | 38 ++++++++++++++--------------- functions/onprem/worker/worker.go | 38 ++++++++++++++--------------- 2 files changed, 38 insertions(+), 38 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 6d0d73ea..7e212505 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -788,42 +788,42 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { } func buildResourcesFromEnv() corev1.ResourceRequirements { - reqs := corev1.ResourceList{} - lims := corev1.ResourceList{} + requests := corev1.ResourceList{} + limits := corev1.ResourceList{} type item struct { - env string - rn corev1.ResourceName - to *corev1.ResourceList + env string + resourceName corev1.ResourceName + resourceList corev1.ResourceList } items := []item{ // kubernetes requests - {env: "SHUFFLE_WORKER_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, - {env: "SHUFFLE_WORKER_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, - {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, + {env: "SHUFFLE_WORKER_CPU_REQUEST", resourceName: corev1.ResourceCPU, resourceList: requests}, + {env: "SHUFFLE_WORKER_MEMORY_REQUEST", resourceName: corev1.ResourceMemory, resourceList: requests}, + {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_REQUEST", resourceName: corev1.ResourceEphemeralStorage, resourceList: requests}, // kubernetes limits - {env: "SHUFFLE_WORKER_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, - {env: "SHUFFLE_WORKER_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, - {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, + {env: "SHUFFLE_WORKER_CPU_LIMIT", resourceName: corev1.ResourceCPU, resourceList: limits}, + {env: "SHUFFLE_WORKER_MEMORY_LIMIT", resourceName: corev1.ResourceMemory, resourceList: limits}, + {env: "SHUFFLE_WORKER_EPHEMERAL_STORAGE_LIMIT", resourceName: corev1.ResourceEphemeralStorage, resourceList: limits}, } for _, it := range items { - if v := strings.TrimSpace(os.Getenv(it.env)); v != "" { - if q, err := resource.ParseQuantity(v); err == nil { - (*it.to)[it.rn] = q + if value := strings.TrimSpace(os.Getenv(it.env)); value != "" { + if quantity, err := resource.ParseQuantity(value); err == nil { + it.resourceList[it.resourceName] = quantity } else { - log.Printf("[WARN] Cannot parse %s=%q as resource quantity: %v", it.env, v, err) + log.Printf("[WARNING] Cannot parse %s=%q as resource quantity: %v", it.env, value, err) } } } rr := corev1.ResourceRequirements{} - if len(reqs) > 0 { - rr.Requests = reqs + if len(requests) > 0 { + rr.Requests = requests } - if len(lims) > 0 { - rr.Limits = lims + if len(limits) > 0 { + rr.Limits = limits } return rr diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 133545f3..3489666c 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -2379,42 +2379,42 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { } func buildResourcesFromEnv() corev1.ResourceRequirements { - reqs := corev1.ResourceList{} - lims := corev1.ResourceList{} + requests := corev1.ResourceList{} + limits := corev1.ResourceList{} type item struct { - env string - rn corev1.ResourceName - to *corev1.ResourceList + env string + resourceName corev1.ResourceName + resourceList corev1.ResourceList } items := []item{ // kubernetes requests - {env: "SHUFFLE_APP_CPU_REQUEST", rn: corev1.ResourceCPU, to: &reqs}, - {env: "SHUFFLE_APP_MEMORY_REQUEST", rn: corev1.ResourceMemory, to: &reqs}, - {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_REQUEST", rn: corev1.ResourceEphemeralStorage, to: &reqs}, + {env: "SHUFFLE_APP_CPU_REQUEST", resourceName: corev1.ResourceCPU, resourceList: requests}, + {env: "SHUFFLE_APP_MEMORY_REQUEST", resourceName: corev1.ResourceMemory, resourceList: requests}, + {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_REQUEST", resourceName: corev1.ResourceEphemeralStorage, resourceList: requests}, // kubernetes limits - {env: "SHUFFLE_APP_CPU_LIMIT", rn: corev1.ResourceCPU, to: &lims}, - {env: "SHUFFLE_APP_MEMORY_LIMIT", rn: corev1.ResourceMemory, to: &lims}, - {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_LIMIT", rn: corev1.ResourceEphemeralStorage, to: &lims}, + {env: "SHUFFLE_APP_CPU_LIMIT", resourceName: corev1.ResourceCPU, resourceList: limits}, + {env: "SHUFFLE_APP_MEMORY_LIMIT", resourceName: corev1.ResourceMemory, resourceList: limits}, + {env: "SHUFFLE_APP_EPHEMERAL_STORAGE_LIMIT", resourceName: corev1.ResourceEphemeralStorage, resourceList: limits}, } for _, it := range items { - if v := strings.TrimSpace(os.Getenv(it.env)); v != "" { - if q, err := resource.ParseQuantity(v); err == nil { - (*it.to)[it.rn] = q + if value := strings.TrimSpace(os.Getenv(it.env)); value != "" { + if quantity, err := resource.ParseQuantity(value); err == nil { + it.resourceList[it.resourceName] = quantity } else { - log.Printf("[WARN] Cannot parse %s=%q as resource quantity: %v", it.env, v, err) + log.Printf("[WARNING] Cannot parse %s=%q as resource quantity: %v", it.env, value, err) } } } rr := corev1.ResourceRequirements{} - if len(reqs) > 0 { - rr.Requests = reqs + if len(requests) > 0 { + rr.Requests = requests } - if len(lims) > 0 { - rr.Limits = lims + if len(limits) > 0 { + rr.Limits = limits } return rr