From 5da98c1d0cace5367d6ca02b4f5744038c686b9e Mon Sep 17 00:00:00 2001 From: Frikky Date: Tue, 16 Sep 2025 16:00:47 +0200 Subject: [PATCH] Re-added SHUFFLE_BASE_IMAGE_NAME usage, and made REGISTRY+BASE IMAGE setting ensure that PullAlways is used to always use latest local image --- functions/onprem/orborus/orborus.go | 23 +++++++++++++++++------ functions/onprem/worker/worker.go | 8 ++++++++ 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 51ab7a69..ab9b2230 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -1058,6 +1058,13 @@ func deployK8sWorker(image string, identifier string, env []string) error { env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_REGISTRY=%s", os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY"))) } + if len(os.Getenv("SHUFFLE_BASE_IMAGE_NAME")) > 0 { + env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", os.Getenv("SHUFFLE_BASE_IMAGE_NAME"))) + } else { + log.Printf("[INFO] SHUFFLE_BASE_IMAGE_NAME is not set. Defaulting to %s", baseimagename) + env = append(env, fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", baseimagename)) + } + if len(os.Getenv("REGISTRY_URL")) > 0 { env = append(env, fmt.Sprintf("REGISTRY_URL=%s", os.Getenv("REGISTRY_URL"))) } @@ -1204,10 +1211,13 @@ func deployK8sWorker(image string, identifier string, env []string) error { ImagePullPolicy: corev1.PullIfNotPresent, } + if len(os.Getenv("REGISTRY_URL")) > 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_NAME")) > 0 { + log.Printf("[INFO] Setting image pull policy to Always as private registry is used.") + containerAttachment.ImagePullPolicy = corev1.PullAlways + } + podname := shuffle.GetPodName() - ctx := context.Background() - if len(podname) > 0 { _, err := shuffle.GetCurrentPodNetworkConfig(ctx, clientset, kubernetesNamespace, podname) if err != nil { @@ -3974,10 +3984,11 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest, image string, identifier := "shuffle-workers" if isKubernetes == "true" { - if shuffle.IsRunningInCluster() { - log.Printf("[INFO] Running in Kubernetes cluster") - // try getting the k8s worker server url - } + // FIXME: Do we need this to map the cluster? + //if shuffle.IsRunningInCluster() { + //log.Printf("[INFO] Running in Kubernetes cluster") + // try getting the k8s worker server url + //} } if strings.Contains(streamUrl, "shuffler.io") || strings.Contains(streamUrl, "localhost") || strings.Contains(streamUrl, "127.0.0.1") || strings.Contains(streamUrl, "shuffle-backend") { diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index fc54c0a1..6dd41de4 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -662,6 +662,14 @@ func deployk8sApp(image string, identifier string, env []string) error { }, } + if len(os.Getenv("REGISTRY_URL")) > 0 && len(os.Getenv("SHUFFLE_BASE_IMAGE_NAME")) > 0 { + log.Printf("[INFO] Setting image pull policy to Always as private registry is used.") + //containerAttachment.ImagePullPolicy = corev1.PullAlways + deployment.Spec.Template.Spec.Containers[0].ImagePullPolicy = corev1.PullAlways + } else { + deployment.Spec.Template.Spec.Containers[0].ImagePullPolicy = corev1.PullIfNotPresent + } + _, err = clientset.AppsV1().Deployments(kubernetesNamespace).Create(context.Background(), deployment, metav1.CreateOptions{}) if err != nil { log.Printf("[ERROR] Failed creating deployment: %v", err)