diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index d1a09ff7..847cd29e 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -49,11 +49,11 @@ import ( //"github.com/mackerelio/go-osstat/memory" //"github.com/shirou/gopsutil/cpu" + appsv1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" rbacv1 "k8s.io/api/rbac/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" - ) // Starts jobs in bulk, so this could be increased @@ -75,7 +75,6 @@ var isKubernetes = os.Getenv("IS_KUBERNETES") var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") var maxCPUPercent = 90 - // var baseimagename = "docker.pkg.github.com/shuffle/shuffle" // var baseimagename = "ghcr.io/frikky" // var baseimagename = "shuffle/shuffle" @@ -104,7 +103,7 @@ var memcached = os.Getenv("SHUFFLE_MEMCACHED") var tenzirUrl = os.Getenv("SHUFFLE_TENZIR_URL") var executionIds = []string{} -var namespacemade = false // For K8s +var namespacemade = false // For K8s var dockercli *dockerclient.Client var containerId string @@ -177,9 +176,26 @@ func getThisContainerId() { log.Printf(`[INFO] Started with containerId "%s"`, containerId) } + +func skipCheckInCleanup(name string) bool { + return strings.HasPrefix(name, "backend") || + strings.HasPrefix(name, "shuffle-backend") || + strings.HasPrefix(name, "frontend") || + strings.HasPrefix(name, "shuffle-frontend") || + strings.HasPrefix(name, "orborus") || + strings.HasPrefix(name, "shuffle-orborus") || + strings.HasPrefix(name, "opensearch") || + strings.HasPrefix(name, "shuffle-opensearch") || + strings.HasPrefix(name, "memcached") || + strings.HasPrefix(name, "shuffle-memcached") +} + func cleanupExistingNodes(ctx context.Context) error { if isKubernetes == "true" { + // of course, this doesn't clean up "nodes" but + // rather pods, services, roles etc. + if kubernetesNamespace == "" { kubernetesNamespace = "default" } @@ -198,6 +214,12 @@ func cleanupExistingNodes(ctx context.Context) error { } for _, pod := range pods.Items { + // check if pod.Name starts with: + // "backend-", "frontend-", "orborus-", "opensearch-" or "memcached-" + if skipCheckInCleanup(pod.Name) { + continue + } + err := clientset.CoreV1().Pods(kubernetesNamespace).Delete(context.Background(), pod.Name, metav1.DeleteOptions{}) if err != nil { log.Printf("[ERROR] Failed deleting pod %s: %s", pod.Name, err) @@ -212,17 +234,40 @@ func cleanupExistingNodes(ctx context.Context) error { } for _, service := range services.Items { + if skipCheckInCleanup(service.Name) { + continue + } + err := clientset.CoreV1().Services(kubernetesNamespace).Delete(context.Background(), service.Name, metav1.DeleteOptions{}) if err != nil { log.Printf("[ERROR] Failed deleting service %s: %s", service.Name, err) } } - log.Printf("[INFO] Cleaned up all pods and services in namespace %s", kubernetesNamespace) + deployments, err := clientset.AppsV1().Deployments(kubernetesNamespace).List(context.Background(), metav1.ListOptions{}) + if err != nil { + log.Printf("[ERROR] Failed listing deployments: %s", err) + return err + } + + for _, deployment := range deployments.Items { + if skipCheckInCleanup(deployment.Name) { + continue + } + + err := clientset.AppsV1().Deployments(kubernetesNamespace).Delete(context.Background(), deployment.Name, metav1.DeleteOptions{}) + if err != nil { + log.Printf("[ERROR] Failed deleting deployment %s: %s", deployment.Name, err) + } + } + + log.Printf("[INFO] Cleaned up all pods and services in namespace %s. Waiting 10 seconds for cleanup to reflect", kubernetesNamespace) + + time.Sleep(10 * time.Second) + return nil } - serviceListOptions := types.ServiceListOptions{} services, err := dockercli.ServiceList( context.Background(), @@ -672,11 +717,10 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { return envVars } - func handleBackendImageDownload(ctx context.Context, images string) error { // Should use docker to: - // 1. Pull the image & tag it - // 2. Distribute the image by updating service if "run" + // 1. Pull the image & tag it + // 2. Distribute the image by updating service if "run" if swarmConfig == "run" || swarmConfig == "swarm" { log.Printf("[DEBUG] Should update service with new image after updating(s): %s. \n\nNOT IMPLEMENTED: Contact support@shuffler.io for support.\n\n", images) @@ -687,8 +731,7 @@ func handleBackendImageDownload(ctx context.Context, images string) error { log.Printf("[DEBUG] Should remove existing image (s): %s", images) // Remove the image - removeOptions := image.RemoveOptions{ - } + removeOptions := image.RemoveOptions{} for _, image := range strings.Split(images, ",") { image = strings.TrimSpace(image) @@ -708,6 +751,120 @@ func handleBackendImageDownload(ctx context.Context, images string) error { return nil } +func fixk8sRoles() { + clientset, _, err := shuffle.GetKubernetesClient() + if err != nil { + log.Printf("[ERROR] Error getting kubernetes client: %s", err) + os.Exit(1) + } + + kubernetesNamespace := "default" + + // Check if namespace exist as variable. If so, make it + if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 { + kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") + } + + // fix roles + // check if "service-creator" role is assigned to the service account "default" + // roleBindingNames := []string{"service-creator-binding", "pod-creator-binding", "deployment-creator-binding"} + serviceAccountName := "default" + roleBindingName := "creator-all" + + resourceTypes := []string{"services", "pods", "deployments"} + + // Check if the RoleBinding exists + roleBinding, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Get(context.TODO(), roleBindingName, metav1.GetOptions{}) + if err != nil { + log.Printf("[WARNING] Failed to get RoleBinding %s: %s", roleBindingName, err) + + // create role and rolebinding + role := &rbacv1.Role{ + ObjectMeta: metav1.ObjectMeta{ + Name: roleBindingName, + }, + Rules: []rbacv1.PolicyRule{ + { + APIGroups: []string{"", "apps"}, + Resources: resourceTypes, + Verbs: []string{"create", "list"}, + }, + }, + } + + ctx := context.TODO() + + _, err := clientset.RbacV1().Roles(kubernetesNamespace).Create(ctx, role, metav1.CreateOptions{}) + if err != nil { + log.Printf("[ERROR] Failed to create Role %s: %s", roleBindingName, err) + if !strings.Contains(fmt.Sprintf("%s", err), "already exists") { + log.Printf("[INFO] role %s already exists", roleBindingName) + } + } + + roleBinding := &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{ + Name: roleBindingName, + }, + Subjects: []rbacv1.Subject{ + { + Kind: "ServiceAccount", + Name: serviceAccountName, + Namespace: kubernetesNamespace, + }, + }, + RoleRef: rbacv1.RoleRef{ + Kind: "Role", + Name: roleBindingName, + }, + } + + _, err = clientset.RbacV1().RoleBindings(kubernetesNamespace).Create(ctx, roleBinding, metav1.CreateOptions{}) + if err != nil { + log.Printf("[ERROR] Failed to create RoleBinding %s: %s", roleBindingName, err) + if strings.Contains(fmt.Sprintf("%s", err), "already exists") { + log.Printf("[INFO] rolebinding %s already exists", roleBindingName) + } + } + + log.Printf("[INFO] Created Role %s and RoleBinding %s", roleBindingName, roleBindingName) + } else { + log.Printf("[INFO] RoleBinding %s exists", roleBindingName) + } + + // Check if the RoleBinding is assigned to the service account + var found bool + for _, subject := range roleBinding.Subjects { + if subject.Kind == "ServiceAccount" && subject.Name == serviceAccountName { + found = true + break + } + } + + if !found { + log.Printf("[WARNING] Service account %s is not assigned to RoleBinding %s\n", serviceAccountName, roleBindingName) + // assign the service account to the rolebinding + roleBinding.Subjects = append(roleBinding.Subjects, rbacv1.Subject{ + Kind: "ServiceAccount", + Name: serviceAccountName, + Namespace: kubernetesNamespace, + }) + + ctx := context.TODO() + + _, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Update(ctx, roleBinding, metav1.UpdateOptions{}) + if err != nil { + log.Printf("[ERROR](ns - %s) Failed to update RoleBinding %s: %s", kubernetesNamespace, roleBindingName, err) + if !strings.Contains(fmt.Sprintf("%s", err), "already exists") { + log.Printf("[INFO] rolebinding %s already exists", roleBindingName) + } + } + } +} + + +func int32Ptr(i int32) *int32 { return &i } + 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"))) @@ -733,7 +890,7 @@ func deployK8sWorker(image string, identifier string, env []string) error { //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 - // The volume mount location is: + // The volume mount location is: // /var/run/secrets/kubernetes.io/serviceaccount // Look for if there is a default service account in use @@ -744,7 +901,6 @@ func deployK8sWorker(image string, identifier string, env []string) error { // 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") @@ -770,9 +926,10 @@ func deployK8sWorker(image string, identifier string, env []string) error { env = append(env, fmt.Sprintf("BASE_URL=%s", baseUrl)) env = append(env, fmt.Sprintf("SHUFFLE_SWARM_CONFIG=%s", swarmConfig)) + env = append(env, fmt.Sprintf("WORKER_HOSTNAME=%s", "shuffle-workers")) if len(kubernetesNamespace) == 0 { - foundNamespace, err := shuffle.GetKubernetesNamespace() + foundNamespace, err := shuffle.GetKubernetesNamespace() if err != nil { //log.Printf("[ERROR] Failed getting Kubernetes namespace: %s", err) } @@ -810,7 +967,7 @@ func deployK8sWorker(image string, identifier string, env []string) error { Name: identifier, Image: kubernetesImage, Env: buildEnvVars(envMap), - + //ImagePullPolicy: "Never", ImagePullPolicy: corev1.PullIfNotPresent, } @@ -829,66 +986,131 @@ func deployK8sWorker(image string, identifier string, env []string) error { } } - // While testing: // kubectl delete pods --all --all-namespaces; kubectl delete services --all --all-namespaces - pod := &corev1.Pod{ + // pod := &corev1.Pod{ + // ObjectMeta: metav1.ObjectMeta{ + // Name: identifier, + // Labels: containerLabels, + // }, + // Spec: corev1.PodSpec{ + // RestartPolicy: "Never", + // // DNSPolicy: "Default", + // DNSPolicy: corev1.DNSClusterFirst, + // // NodeSelector: map[string]string{ + // // "node": "master", + // // }, + // Containers: []corev1.Container{ + // containerAttachment, + // }, + // }, + // } + + // // Check if running on ARM or x86 to download the correct image + + // // Get current pod's network so we can make the pod in it + + // _, 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 { + // //log.Printf("[ERROR] Failed creating pod: %v", err) + // return err + // } + + // log.Printf("[INFO] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace) + + // // kubectl expose pod shuffle-workers --type=LoadBalancer --port=33333 + // service := &corev1.Service{ + // ObjectMeta: metav1.ObjectMeta{ + // Name: identifier, + // }, + // Spec: corev1.ServiceSpec{ + // Selector: map[string]string{ + // "container": "shuffle-workers", + // }, + // Ports: []corev1.ServicePort{ + // { + // Protocol: "TCP", + // Port: 33333, + // TargetPort: intstr.FromInt(33333), + // }, + // }, + // Type: corev1.ServiceTypeLoadBalancer, + // }, + // } + + // _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + // if err != nil { + // log.Printf("[ERROR] Failed creating service: %v", err) + // return err + // } + + 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 { + replicaNumber = tmpInt + + } + } + + replicaNumberInt32 := int32(replicaNumber) + + deployment := &appsv1.Deployment{ ObjectMeta: metav1.ObjectMeta{ - Name: identifier, - Labels: containerLabels, + Name: identifier, }, - Spec: corev1.PodSpec{ - RestartPolicy: "Never", - // DNSPolicy: "Default", - DNSPolicy: corev1.DNSClusterFirst, - // NodeSelector: map[string]string{ - // "node": "master", - // }, - Containers: []corev1.Container{ - containerAttachment, + Spec: appsv1.DeploymentSpec{ + Replicas: int32Ptr(replicaNumberInt32), + Selector: &metav1.LabelSelector{ + MatchLabels: containerLabels, + }, + Template: corev1.PodTemplateSpec{ + ObjectMeta: metav1.ObjectMeta{ + Labels: containerLabels, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + containerAttachment, + }, + DNSPolicy: corev1.DNSClusterFirst, + }, }, }, } - // Check if running on ARM or x86 to download the correct image - - // Get current pod's network so we can make the pod in it - - _, err = clientset.CoreV1().Pods(kubernetesNamespace).List(context.Background(), metav1.ListOptions{}) + _, err = clientset.AppsV1().Deployments(kubernetesNamespace).Create(context.Background(), deployment, metav1.CreateOptions{}) 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 { - //log.Printf("[ERROR] Failed creating pod: %v", err) + log.Printf("[ERROR] Failed creating deployment: %v", err) return err } - log.Printf("[INFO] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace) - - // kubectl expose pod shuffle-workers --type=LoadBalancer --port=33333 + // kubectl expose deployment shuffle-workers --type=NodePort --port=33333 --target-port=33333 service := &corev1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: identifier, }, Spec: corev1.ServiceSpec{ - Selector: map[string]string{ - "container": "shuffle-workers", - }, + Selector: containerLabels, Ports: []corev1.ServicePort{ { - Protocol: "TCP", - Port: 33333, + Protocol: "TCP", + Port: 33333, TargetPort: intstr.FromInt(33333), }, }, - Type: corev1.ServiceTypeLoadBalancer, + Type: corev1.ServiceTypeNodePort, }, } - _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.TODO(), service, metav1.CreateOptions{}) + _, err = clientset.CoreV1().Services(kubernetesNamespace).Create(context.Background(), service, metav1.CreateOptions{}) if err != nil { log.Printf("[ERROR] Failed creating service: %v", err) return err @@ -897,7 +1119,6 @@ func deployK8sWorker(image string, identifier string, env []string) error { return nil } - func deployWorker(image string, identifier string, env []string, executionRequest shuffle.ExecutionRequest) error { if len(os.Getenv("REGISTRY_URL")) > 0 && os.Getenv("REGISTRY_URL") != "" { env = append(env, fmt.Sprintf("REGISTRY_URL=%s", os.Getenv("REGISTRY_URL"))) @@ -1299,8 +1520,7 @@ func getOrborusStats(ctx context.Context) shuffle.OrborusStats { newStats.MaxMemory = int(pers.MemTotal) } - - // Get list of all running containers + // Get list of all running containers containers, err := dockercli.ContainerList(ctx, container.ListOptions{}) if err != nil { @@ -1414,10 +1634,6 @@ func getOrborusStats(ctx context.Context) shuffle.OrborusStats { return newStats } - - - - func sendRemoveRequest(client *http.Client, toBeRemoved shuffle.ExecutionRequestWrapper, baseUrl, environment, auth, org string, sleepTime int) error { confirmUrl := fmt.Sprintf("%s/api/v1/workflows/queue/confirm", baseUrl) @@ -1496,112 +1712,7 @@ func main() { } if isKubernetes == "true" { - clientset, _, err := shuffle.GetKubernetesClient() - if err != nil { - log.Printf("[ERROR] Error getting kubernetes client: %s", err) - os.Exit(1) - } - - kubernetesNamespace := "default" - - // Check if namespace exist as variable. If so, make it - if len(os.Getenv("KUBERNETES_NAMESPACE")) > 0 && !namespacemade { - kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE") - } - - // fix roles - // check if "service-creator" role is assigned to the service account "default" - roleBindingName := "service-creator-binding" - serviceAccountName := "default" - // Check if the RoleBinding exists - roleBinding, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Get(context.TODO(), roleBindingName, metav1.GetOptions{}) - if err != nil { - log.Printf("[WARNING] Failed to get RoleBinding %s: %s", roleBindingName, err) - // create role and rolebinding - role := &rbacv1.Role{ - ObjectMeta: metav1.ObjectMeta{ - Name: roleBindingName, - }, - Rules: []rbacv1.PolicyRule{ - { - APIGroups: []string{""}, - Resources: []string{"services"}, - Verbs: []string{"create"}, - }, - }, - } - - ctx := context.TODO() - - _, err := clientset.RbacV1().Roles(kubernetesNamespace).Create(ctx, role, metav1.CreateOptions{}) - if err != nil { - log.Printf("[ERROR] Failed to create Role %s: %s", roleBindingName, err) - if !strings.Contains(fmt.Sprintf("%s", err), "already exists") { - log.Printf("[INFO] role %s already exists", roleBindingName) - } - } - - roleBinding := &rbacv1.RoleBinding{ - ObjectMeta: metav1.ObjectMeta{ - Name: roleBindingName, - }, - Subjects: []rbacv1.Subject{ - { - Kind: "ServiceAccount", - Name: serviceAccountName, - Namespace: kubernetesNamespace, - }, - }, - RoleRef: rbacv1.RoleRef{ - Kind: "Role", - Name: roleBindingName, - }, - } - - - - _, err = clientset.RbacV1().RoleBindings(kubernetesNamespace).Create(ctx, roleBinding, metav1.CreateOptions{}) - if err != nil { - log.Printf("[ERROR] Failed to create RoleBinding %s: %s", roleBindingName, err) - if !strings.Contains(fmt.Sprintf("%s", err), "already exists") { - log.Printf("[INFO] rolebinding %s already exists", roleBindingName) - } - } - - - log.Printf("[INFO] Created Role %s and RoleBinding %s", roleBindingName, roleBindingName) - } else { - log.Printf("[INFO] RoleBinding %s exists", roleBindingName) - } - - // Check if the RoleBinding is assigned to the service account - var found bool - for _, subject := range roleBinding.Subjects { - if subject.Kind == "ServiceAccount" && subject.Name == serviceAccountName { - found = true - break - } - } - - if !found { - log.Printf("[WARNING] Service account %s is not assigned to RoleBinding %s\n", serviceAccountName, roleBindingName) - // assign the service account to the rolebinding - roleBinding.Subjects = append(roleBinding.Subjects, rbacv1.Subject{ - Kind: "ServiceAccount", - Name: serviceAccountName, - Namespace: kubernetesNamespace, - }) - - ctx := context.TODO() - - _, err := clientset.RbacV1().RoleBindings(kubernetesNamespace).Update(ctx, roleBinding, metav1.UpdateOptions{}) - if err != nil { - log.Printf("[ERROR] Failed to update RoleBinding %s: %s", roleBindingName, err) - if !strings.Contains(fmt.Sprintf("%s", err), "already exists") { - log.Printf("[INFO] rolebinding %s already exists", roleBindingName) - } - } - } + fixk8sRoles() } startupDelay := os.Getenv("SHUFFLE_ORBORUS_STARTUP_DELAY") @@ -1711,7 +1822,7 @@ func main() { log.Printf("[DEBUG] Cleaning up containers from previous run") cleanupExistingNodes(ctx) - time.Sleep(time.Duration(5) * time.Second) + time.Sleep(time.Duration(5) * time.Second) log.Printf("[DEBUG] Deploying worker image %s to swarm", workerImage) @@ -1879,13 +1990,12 @@ func main() { continue } - if hasStarted && len(executionRequests.Data) > 0 { //log.Printf("[INFO] Body: %s", string(body)) // Type string `json:"type"` } - // FIXME: Add features here for orborus & worker to + // FIXME: Add features here for orborus & worker to // do things on behalf of backend var toBeRemoved shuffle.ExecutionRequestWrapper if len(executionRequests.Data) > 0 { @@ -1957,7 +2067,7 @@ func main() { log.Printf("[WARNING] Throttle - Cutting down requests from %d to %d (MAX: %d, CUR: %d)", len(executionRequests.Data), allowed, maxConcurrency, executionCount) executionRequests.Data = executionRequests.Data[0:allowed] } - } else if (swarmControlMode && (swarmConfig == "run" || swarmConfig == "swarm")) { + } else if swarmControlMode && (swarmConfig == "run" || swarmConfig == "swarm") { if len(executionRequests.Data) > 50 { executionRequests.Data = executionRequests.Data[0:50] } @@ -2106,7 +2216,6 @@ func main() { } } - // func deployPipeline(image, identifier, command string) error { // if isKubernetes == "true" { // return errors.New("Kubernetes not implemented") @@ -2131,7 +2240,6 @@ func main() { // envVariables := []string{ // } - // // Add volume binds for storage // // Want read/write with full access for the container // //sourceFolder := "/Users/frikky/git/shuffle/shuffle-database" @@ -2168,8 +2276,7 @@ func main() { // config.Labels = map[string]string{ // "name": identifier, // "shuffle": "shuffle", -// } - +// } // cont, err := dockercli.ContainerCreate( // ctx, @@ -2191,8 +2298,8 @@ func main() { // containerStartOptions := container.StartOptions{} // err = dockercli.ContainerStart( -// ctx, -// cont.ID, +// ctx, +// cont.ID, // containerStartOptions, // ) // if err != nil { @@ -2211,8 +2318,8 @@ func main() { // } // err = dockercli.ContainerStart( -// ctx, -// cont.ID, +// ctx, +// cont.ID, // containerStartOptions, // ) // if err != nil { @@ -2256,8 +2363,6 @@ func main() { // return nil // } - - // Tenzir command samples // docker pull ghcr.io/dominiklohmann/tenzir-arm64:latest // docker tag ghcr.io/dominiklohmann/tenzir-arm64:latest tenzir/tenzir:latest @@ -2265,14 +2370,14 @@ func main() { // Read from Cache and send it to a webhook // docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines' func handlePipeline(incRequest shuffle.ExecutionRequest) error { - + if tenzirUrl == "" { tenzirUrl = "http://localhost:5160" - log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s",tenzirUrl) + log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s", tenzirUrl) } err := deployTenzirNode() - if err != nil{ + if err != nil { log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err) return err } @@ -2306,7 +2411,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { if err != nil { log.Printf("[ERROR] Failed Deleting Pipeline %s", err) return err - } + } } else if incRequest.Type == "PIPELINE_STOP" { log.Printf("[INFO] Should stop the pipeline %#v", identifier) pipelineId, err := searchPipeline(identifier) @@ -2322,10 +2427,10 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { log.Printf("[INFO] successfully stopped the Pipeline: %s", pipelineId) } - } else if incRequest.Type == "PIPELINE_START" { + } else if incRequest.Type == "PIPELINE_START" { log.Printf("[INFO] Should start the pipeline %#v", identifier) pipelineId, err := searchPipeline(identifier) - if err != nil { + if err != nil { if err.Error() == "no existing pipeline found with name" { log.Printf("[WARNING] no pipeline found for %s, creating a new one", identifier) _, CreateErr := createPipeline(command, identifier) @@ -2351,157 +2456,157 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { } func deployTenzirNode() error { - if isKubernetes == "true" { - return errors.New("kubernetes not implemented") - } + if isKubernetes == "true" { + return errors.New("kubernetes not implemented") + } - ctx := context.Background() - cacheKey := "tenzir-key" + ctx := context.Background() + cacheKey := "tenzir-key" - imageName := "tenzir/tenzir:latest" - containerName := "tenzir-node" - containerStartOptions := container.StartOptions{} + imageName := "tenzir/tenzir:latest" + containerName := "tenzir-node" + containerStartOptions := container.StartOptions{} - _, err := shuffle.GetCache(ctx, cacheKey) - if err == nil { - return nil - } + _, err := shuffle.GetCache(ctx, cacheKey) + if err == nil { + return nil + } - containerInfo, err := dockercli.ContainerInspect(ctx, containerName) - if err != nil { - if dockerclient.IsErrNotFound(err) { + containerInfo, err := dockercli.ContainerInspect(ctx, containerName) + if err != nil { + if dockerclient.IsErrNotFound(err) { - // Check if image exists - _, _, err := dockercli.ImageInspectWithRaw(ctx, imageName) - if dockerclient.IsErrNotFound(err) { - log.Printf("[DEBUG] pulling image %s", imageName) - pullOptions := image.PullOptions{} - out, err := dockercli.ImagePull(ctx, imageName, pullOptions) - if err != nil { - log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err) - return err - } - defer out.Close() + // Check if image exists + _, _, err := dockercli.ImageInspectWithRaw(ctx, imageName) + if dockerclient.IsErrNotFound(err) { + log.Printf("[DEBUG] pulling image %s", imageName) + pullOptions := image.PullOptions{} + out, err := dockercli.ImagePull(ctx, imageName, pullOptions) + if err != nil { + log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err) + return err + } + defer out.Close() - io.Copy(io.Discard, out) - } else if err != nil { - return err - } + io.Copy(io.Discard, out) + } else if err != nil { + return err + } - err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions) - if err != nil { - return err - } - } else { - return err - } - } else { - if !containerInfo.State.Running { - log.Printf("[DEBUG] Tenzir Node exists but is not running") - err := dockercli.ContainerStart(ctx, containerName, containerStartOptions) - if err != nil { - log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) - return err - } + err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions) + if err != nil { + return err + } + } else { + return err + } + } else { + if !containerInfo.State.Running { + log.Printf("[DEBUG] Tenzir Node exists but is not running") + err := dockercli.ContainerStart(ctx, containerName, containerStartOptions) + if err != nil { + log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) + return err + } - log.Printf("[INFO] Waiting for Tenzir to become available ...") - err = checkTenzirNode() - if err != nil { - return err - } - } - } + log.Printf("[INFO] Waiting for Tenzir to become available ...") + err = checkTenzirNode() + if err != nil { + return err + } + } + } - tenzirStatus := struct { - ContainerStatus string `json:"container_status"` - }{ - ContainerStatus: "running", - } + tenzirStatus := struct { + ContainerStatus string `json:"container_status"` + }{ + ContainerStatus: "running", + } - cacheData, err := json.Marshal(tenzirStatus) - if err != nil { - log.Printf("[WARNING] Failed marshalling execution: %s", err) - } - err = shuffle.SetCache(ctx, cacheKey, cacheData, 1) - if err != nil { - log.Printf("[WARNING] Failed updating cache for tenzir: %s", err) - } + cacheData, err := json.Marshal(tenzirStatus) + if err != nil { + log.Printf("[WARNING] Failed marshalling execution: %s", err) + } + err = shuffle.SetCache(ctx, cacheKey, cacheData, 1) + if err != nil { + log.Printf("[WARNING] Failed updating cache for tenzir: %s", err) + } - return nil + return nil } func checkTenzirNode() error { - retries := 20 - retryInterval := 3 * time.Second - url := fmt.Sprintf("%s/api/v0/ping",tenzirUrl) + retries := 20 + retryInterval := 3 * time.Second + url := fmt.Sprintf("%s/api/v0/ping", tenzirUrl) forwardMethod := "POST" - client := http.Client{} - req, err := http.NewRequest(forwardMethod, url, nil) + client := http.Client{} + req, err := http.NewRequest(forwardMethod, url, nil) if err != nil { log.Printf("[ERROR] Failed to create HTTP request: %s", err) return err } - for i := 0; i < retries; i++ { - resp, err := client.Do(req) - if err == nil && resp.StatusCode == http.StatusOK { - return nil - } - time.Sleep(retryInterval) - } + for i := 0; i < retries; i++ { + resp, err := client.Do(req) + if err == nil && resp.StatusCode == http.StatusOK { + return nil + } + time.Sleep(retryInterval) + } - return fmt.Errorf("tenzir node is not available") + return fmt.Errorf("tenzir node is not available") } func createAndStartTenzirNode(ctx context.Context, containerName, imageName string, containerStartOptions container.StartOptions) error { - healthconfig := &container.HealthConfig{ - Test: []string{"tenzir --connection-timeout=30s --connection-retry-delay=1s 'api /ping'"}, - Interval: 30 * time.Second, - Retries: 1, - } + healthconfig := &container.HealthConfig{ + Test: []string{"tenzir --connection-timeout=30s --connection-retry-delay=1s 'api /ping'"}, + Interval: 30 * time.Second, + Retries: 1, + } - config := &container.Config{ - Cmd: []string{"--commands=web server --mode=dev --bind=0.0.0.0"}, - Image: imageName, - Healthcheck: healthconfig, - ExposedPorts: nat.PortSet{"5160/tcp": struct{}{}}, - Entrypoint: []string{containerName}, - } + config := &container.Config{ + Cmd: []string{"--commands=web server --mode=dev --bind=0.0.0.0"}, + Image: imageName, + Healthcheck: healthconfig, + ExposedPorts: nat.PortSet{"5160/tcp": struct{}{}}, + Entrypoint: []string{containerName}, + } - hostConfig := &container.HostConfig{ - PortBindings: nat.PortMap{ - "5160/tcp": []nat.PortBinding{{HostPort: "5160"}}, - }, - Mounts: []mount.Mount{ - { - Type: mount.TypeVolume, - Source: containerName, - Target: "/var/lib/tenzir/", - }, - }, - VolumeDriver: "local", - } - _, err := dockercli.ContainerCreate(ctx, config, hostConfig, nil, nil, containerName) - if err != nil { - return err - } + hostConfig := &container.HostConfig{ + PortBindings: nat.PortMap{ + "5160/tcp": []nat.PortBinding{{HostPort: "5160"}}, + }, + Mounts: []mount.Mount{ + { + Type: mount.TypeVolume, + Source: containerName, + Target: "/var/lib/tenzir/", + }, + }, + VolumeDriver: "local", + } + _, err := dockercli.ContainerCreate(ctx, config, hostConfig, nil, nil, containerName) + if err != nil { + return err + } - err = dockercli.ContainerStart(ctx, containerName, containerStartOptions) - if err != nil { - log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) - return err - } - log.Printf("[INFO] Tenzir Node container started successfully") + err = dockercli.ContainerStart(ctx, containerName, containerStartOptions) + if err != nil { + log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) + return err + } + log.Printf("[INFO] Tenzir Node container started successfully") - log.Printf("[INFO] Waiting for Tenzir to become available ...") - err = checkTenzirNode() - if err != nil { - return err - } - log.Printf("[INFO] Successfully deployed Tenzir Node !") + log.Printf("[INFO] Waiting for Tenzir to become available ...") + err = checkTenzirNode() + if err != nil { + return err + } + log.Printf("[INFO] Successfully deployed Tenzir Node !") - return nil + return nil } func createPipeline(command, identifier string) (string, error) { @@ -2509,7 +2614,7 @@ func createPipeline(command, identifier string) (string, error) { toBeDeleted := false pipelineId, err := searchPipeline(identifier) - url := fmt.Sprintf("%s/api/v0/pipeline/create", tenzirUrl) + url := fmt.Sprintf("%s/api/v0/pipeline/create", tenzirUrl) forwardMethod := "POST" if err != nil { @@ -2522,7 +2627,7 @@ func createPipeline(command, identifier string) (string, error) { log.Printf("[INFO] an existing pipeline found with ID: %s. it will be deleted", pipelineId) toBeDeleted = true } - if strings.Contains(command, "shuffler.io") { + if strings.Contains(command, "shuffler.io") { } else { var scheme string @@ -2536,7 +2641,7 @@ func createPipeline(command, identifier string) (string, error) { if startIndex != -1 { endIndex := startIndex + len(scheme) endIndex += strings.Index(command[endIndex:], "/") - + command = command[:startIndex] + baseUrl + command[endIndex:] } } @@ -2615,7 +2720,7 @@ func createPipeline(command, identifier string) (string, error) { func updatePipelineState(pipelineId, action string) (string, error) { - url := fmt.Sprintf("%s/api/v0/pipeline/update", tenzirUrl) + url := fmt.Sprintf("%s/api/v0/pipeline/update", tenzirUrl) forwardMethod := "POST" requestBody := map[string]interface{}{ @@ -2684,7 +2789,7 @@ func deletePipeline(pipelineId string) error { "id": pipelineId, } - url := fmt.Sprintf("%s/api/v0/pipeline/delete", tenzirUrl) + url := fmt.Sprintf("%s/api/v0/pipeline/delete", tenzirUrl) forwardMethod := "POST" requestBodyJSON, err := json.Marshal(requestBody) @@ -2730,9 +2835,9 @@ func searchPipeline(identifier string) (string, error) { Name string `json:"name"` } - var reqBody []byte + var reqBody []byte - url := fmt.Sprintf("%s/api/v0/pipeline/list", tenzirUrl) + url := fmt.Sprintf("%s/api/v0/pipeline/list", tenzirUrl) resp, err := http.Post(url, "application/json", bytes.NewBuffer(reqBody)) if err != nil { @@ -2817,10 +2922,9 @@ func searchPipeline(identifier string) (string, error) { func getRunningWorkers(ctx context.Context, workerTimeout int) int { //log.Printf("[DEBUG] Getting running workers with API version %s", dockerApiVersion) counter := 0 - if isKubernetes == "true" { + if isKubernetes == "true" { log.Printf("[INFO] Getting running workers in kubernetes") - thresholdTime := time.Now().Add(time.Duration(-workerTimeout) * time.Second) clientset, _, err := shuffle.GetKubernetesClient() @@ -3128,7 +3232,7 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest, image string, debugCommand := fmt.Sprintf("docker service logs shuffle-workers 2>&1 -f | grep %s", workflowExecution.ExecutionId) if isKubernetes == "true" { - debugCommand = fmt.Sprintf("kubectl logs -n %s %s | grep %s", kubernetesNamespace, identifier, workflowExecution.ExecutionId) + debugCommand = fmt.Sprintf("kubectl logs -n %s container=shuffle-worker | grep %s", kubernetesNamespace, workflowExecution.ExecutionId) } log.Printf("[DEBUG] Ran worker from request with execution ID: %s. Worker URL: %s. DEBUGGING:\n%s", workflowExecution.ExecutionId, streamUrl, debugCommand)