diff --git a/functions/k8s-confs/backend-deployment.yaml b/functions/k8s-confs/backend-deployment.yaml new file mode 100644 index 00000000..9dc1b390 --- /dev/null +++ b/functions/k8s-confs/backend-deployment.yaml @@ -0,0 +1,391 @@ +--- + +apiVersion: v1 +kind: PersistentVolume +metadata: + name: shuffle-pv +spec: + capacity: + storage: 10Gi + accessModes: + - ReadWriteOnce + persistentVolumeReclaimPolicy: Retain + storageClassName: shuffle-storage + hostPath: + path: /mnt/shuffle-data/backend + +--- + +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + creationTimestamp: null + labels: + io.kompose.service: backend-files-claim + name: backend-files-claim +spec: + accessModes: + - ReadWriteOnce + resources: + requests: + storage: 1Gi +status: {} + +--- + +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + creationTimestamp: null + labels: + io.kompose.service: backend-apps-claim + name: backend-apps-claim +spec: + accessModes: + - ReadWriteOnce + resources: + requests: + storage: 1Gi +status: {} + +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: backend + name: backend +spec: + replicas: 1 + selector: + matchLabels: + io.kompose.service: backend + strategy: + type: Recreate + template: + metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.network/shuffle: "true" + io.kompose.service: backend + app: shuffle-backend + name: shuffle-backend + spec: + volumes: + - name: shuffle-files + persistentVolumeClaim: + claimName: backend-files-claim + - name: shuffle-apps + persistentVolumeClaim: + claimName: backend-apps-claim + - name: docker-socket + hostPath: + path: /var/run/docker.sock + containers: + - env: + - name: BACKEND_HOSTNAME + valueFrom: + configMapKeyRef: + key: BACKEND_HOSTNAME + name: env + - name: BACKEND_PORT + valueFrom: + configMapKeyRef: + key: BACKEND_PORT + name: env + - name: BASE_URL + valueFrom: + configMapKeyRef: + key: BASE_URL + name: env + - name: DATASTORE_EMULATOR_HOST + valueFrom: + configMapKeyRef: + key: DATASTORE_EMULATOR_HOST + name: env + - name: DB_LOCATION + valueFrom: + configMapKeyRef: + key: DB_LOCATION + name: env + - name: DOCKER_API_VERSION + valueFrom: + configMapKeyRef: + key: DOCKER_API_VERSION + name: env + - name: ENVIRONMENT_NAME + valueFrom: + configMapKeyRef: + key: ENVIRONMENT_NAME + name: env + - name: FRONTEND_PORT + valueFrom: + configMapKeyRef: + key: FRONTEND_PORT + name: env + - name: FRONTEND_PORT_HTTPS + valueFrom: + configMapKeyRef: + key: FRONTEND_PORT_HTTPS + name: env + - name: HTTPS_PROXY + valueFrom: + configMapKeyRef: + key: HTTPS_PROXY + name: env + - name: HTTP_PROXY + valueFrom: + configMapKeyRef: + key: HTTP_PROXY + name: env + - name: ORBORUS_CONTAINER_NAME + valueFrom: + configMapKeyRef: + key: ORBORUS_CONTAINER_NAME + name: env + - name: ORG_ID + valueFrom: + configMapKeyRef: + key: ORG_ID + name: env + - name: OUTER_HOSTNAME + valueFrom: + configMapKeyRef: + key: OUTER_HOSTNAME + name: env + - name: SHUFFLE_APP_DOWNLOAD_LOCATION + valueFrom: + configMapKeyRef: + key: SHUFFLE_APP_DOWNLOAD_LOCATION + name: env + - name: SHUFFLE_APP_FORCE_UPDATE + valueFrom: + configMapKeyRef: + key: SHUFFLE_APP_FORCE_UPDATE + name: env + - name: SHUFFLE_APP_HOTLOAD_FOLDER + valueFrom: + configMapKeyRef: + key: SHUFFLE_APP_HOTLOAD_FOLDER + name: env + - name: SHUFFLE_APP_HOTLOAD_LOCATION + valueFrom: + configMapKeyRef: + key: SHUFFLE_APP_HOTLOAD_LOCATION + name: env + - name: SHUFFLE_BASE_IMAGE_NAME + valueFrom: + configMapKeyRef: + key: SHUFFLE_BASE_IMAGE_NAME + name: env + - name: SHUFFLE_BASE_IMAGE_REGISTRY + valueFrom: + configMapKeyRef: + key: SHUFFLE_BASE_IMAGE_REGISTRY + name: env + - name: SHUFFLE_BASE_IMAGE_TAG_SUFFIX + valueFrom: + configMapKeyRef: + key: SHUFFLE_BASE_IMAGE_TAG_SUFFIX + name: env + - name: SHUFFLE_CHAT_DISABLED + valueFrom: + configMapKeyRef: + key: SHUFFLE_CHAT_DISABLED + name: env + - name: SHUFFLE_CONTAINER_AUTO_CLEANUP + valueFrom: + configMapKeyRef: + key: SHUFFLE_CONTAINER_AUTO_CLEANUP + name: env + - name: SHUFFLE_DEFAULT_APIKEY + valueFrom: + configMapKeyRef: + key: SHUFFLE_DEFAULT_APIKEY + name: env + - name: SHUFFLE_DEFAULT_PASSWORD + valueFrom: + configMapKeyRef: + key: SHUFFLE_DEFAULT_PASSWORD + name: env + - name: SHUFFLE_DEFAULT_USERNAME + valueFrom: + configMapKeyRef: + key: SHUFFLE_DEFAULT_USERNAME + name: env + - name: SHUFFLE_DOWNLOAD_AUTH_BRANCH + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_AUTH_BRANCH + name: env + - name: SHUFFLE_DOWNLOAD_AUTH_PASSWORD + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_AUTH_PASSWORD + name: env + - name: SHUFFLE_DOWNLOAD_AUTH_USERNAME + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_AUTH_USERNAME + name: env + - name: SHUFFLE_DOWNLOAD_WORKFLOW_BRANCH + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_WORKFLOW_BRANCH + name: env + - name: SHUFFLE_DOWNLOAD_WORKFLOW_LOCATION + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_WORKFLOW_LOCATION + name: env + - name: SHUFFLE_DOWNLOAD_WORKFLOW_PASSWORD + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_WORKFLOW_PASSWORD + name: env + - name: SHUFFLE_DOWNLOAD_WORKFLOW_USERNAME + valueFrom: + configMapKeyRef: + key: SHUFFLE_DOWNLOAD_WORKFLOW_USERNAME + name: env + - name: SHUFFLE_ELASTIC + valueFrom: + configMapKeyRef: + key: SHUFFLE_ELASTIC + name: env + - name: SHUFFLE_ENCRYPTION_MODIFIER + valueFrom: + configMapKeyRef: + key: SHUFFLE_ENCRYPTION_MODIFIER + name: env + - name: SHUFFLE_FILE_LOCATION + valueFrom: + configMapKeyRef: + key: SHUFFLE_FILE_LOCATION + name: env + - name: SHUFFLE_LOGS_DISABLED + valueFrom: + configMapKeyRef: + key: SHUFFLE_LOGS_DISABLED + name: env + - name: SHUFFLE_OPENSEARCH_APIKEY + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_APIKEY + name: env + - name: SHUFFLE_OPENSEARCH_CERTIFICATE_FILE + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_CERTIFICATE_FILE + name: env + - name: SHUFFLE_OPENSEARCH_CLOUDID + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_CLOUDID + name: env + - name: SHUFFLE_OPENSEARCH_INDEX_PREFIX + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_INDEX_PREFIX + name: env + - name: SHUFFLE_OPENSEARCH_PASSWORD + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_PASSWORD + name: env + - name: SHUFFLE_OPENSEARCH_PROXY + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_PROXY + name: env + - name: SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY + name: env + - name: SHUFFLE_OPENSEARCH_URL + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_URL + name: env + - name: SHUFFLE_OPENSEARCH_USERNAME + valueFrom: + configMapKeyRef: + key: SHUFFLE_OPENSEARCH_USERNAME + name: env + - name: SHUFFLE_ORBORUS_STARTUP_DELAY + valueFrom: + configMapKeyRef: + key: SHUFFLE_ORBORUS_STARTUP_DELAY + name: env + - name: SHUFFLE_PASS_APP_PROXY + valueFrom: + configMapKeyRef: + key: SHUFFLE_PASS_APP_PROXY + name: env + - name: SHUFFLE_PASS_WORKER_PROXY + valueFrom: + configMapKeyRef: + key: SHUFFLE_PASS_WORKER_PROXY + name: env + - name: SHUFFLE_RERUN_SCHEDULE + valueFrom: + configMapKeyRef: + key: SHUFFLE_RERUN_SCHEDULE + name: env + - name: SSO_REDIRECT_URL + valueFrom: + configMapKeyRef: + key: SSO_REDIRECT_URL + name: env + - name: TZ + valueFrom: + configMapKeyRef: + key: TZ + name: env + image: ghcr.io/shuffle/shuffle-backend:latest + name: shuffle-backend + ports: + - containerPort: 5001 + resources: {} + volumeMounts: + - name: docker-socket + mountPath: /var/run/docker.sock + - name: shuffle-apps + mountPath: /shuffle-apps + - name: shuffle-files + mountPath: /shuffle-files + restartPolicy: Always +status: {} + +--- + +apiVersion: v1 +kind: Service +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: backend + name: shuffle-backend +spec: + type: NodePort + ports: + - name: "50001" + port: 5001 + targetPort: 5001 + nodePort: 30009 + selector: + io.kompose.service: backend +status: + loadBalancer: {} + diff --git a/functions/k8s-confs/frontend-deployment.yaml b/functions/k8s-confs/frontend-deployment.yaml new file mode 100644 index 00000000..cb11f728 --- /dev/null +++ b/functions/k8s-confs/frontend-deployment.yaml @@ -0,0 +1,68 @@ +--- + +apiVersion: apps/v1 +kind: Deployment +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: frontend + name: frontend +spec: + replicas: 1 + selector: + matchLabels: + io.kompose.service: frontend + strategy: {} + template: + metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.network/shuffle: "true" + io.kompose.service: frontend + spec: + containers: + - env: + - name: BACKEND_HOSTNAME + image: ghcr.io/shuffle/shuffle-frontend:latest + name: shuffle-frontend + ports: + - containerPort: 80 + - containerPort: 443 + resources: {} + hostname: shuffle-frontend + restartPolicy: Always +status: {} + +--- + +apiVersion: v1 +kind: Service +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: frontend + name: frontend +spec: + type: NodePort + ports: + - name: "80" + port: 80 + targetPort: 80 + nodePort: 30007 + - name: "443" + port: 443 + targetPort: 443 + nodePort: 30008 + selector: + io.kompose.service: frontend +# status: +# loadBalancer: {} diff --git a/functions/k8s-confs/opensearch-deployment.yaml b/functions/k8s-confs/opensearch-deployment.yaml new file mode 100644 index 00000000..37e2e71a --- /dev/null +++ b/functions/k8s-confs/opensearch-deployment.yaml @@ -0,0 +1,123 @@ +--- + +apiVersion: v1 +kind: PersistentVolume +metadata: + name: shuffle-os-pv +spec: + capacity: + storage: 10Gi # Adjust the storage size as per your requirements + accessModes: + - ReadWriteOnce # This allows read-write access to a single node + persistentVolumeReclaimPolicy: Retain # Adjust the reclaim policy as per your needs + storageClassName: shuffle-storage # Set the desired storage class + hostPath: + path: /mnt/shuffle-data/open-search + +--- + +apiVersion: v1 +kind: PersistentVolumeClaim +metadata: + creationTimestamp: null + labels: + io.kompose.service: opensearch-claim0 + name: opensearch-claim0 +spec: + accessModes: + - ReadWriteOnce + resources: + requests: + storage: 500Mi +status: {} + +--- +apiVersion: apps/v1 +kind: Deployment +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: opensearch + name: opensearch +spec: + replicas: 1 + selector: + matchLabels: + io.kompose.service: opensearch + strategy: {} + template: + metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.network/shuffle: "true" + io.kompose.service: opensearch + spec: + containers: + - env: + - name: OPENSEARCH_JAVA_OPTS + value: -Xms1024m -Xmx1024m + #- name: bootstrap.memory_lock + #value: "true" + - name: cluster.initial_master_nodes + value: shuffle-opensearch + - name: cluster.name + value: shuffle-cluster + - name: cluster.routing.allocation.disk.threshold_enabled + value: "false" + - name: discovery.seed_hosts + value: shuffle-opensearch + - name: node.name + value: shuffle-opensearch + - name: node.store.allow_mmap + value: "false" + - name: DB_LOCATION + valueFrom: + configMapKeyRef: + name: env + key: DB_LOCATION + image: opensearchproject/opensearch:2.5.0 + name: shuffle-opensearch + ports: + - containerPort: 9200 + resources: {} + volumeMounts: + - mountPath: /usr/share/opensearch/data + name: opensearch-claim0 + hostname: shuffle-opensearch + restartPolicy: Always + volumes: + - name: opensearch-claim0 + persistentVolumeClaim: + claimName: opensearch-claim0 +status: {} + +--- + +apiVersion: v1 +kind: Service +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: opensearch + name: opensearch +spec: + type: NodePort + ports: + - name: "9200" + port: 9200 + targetPort: 9200 + nodePort: 31001 + selector: + io.kompose.service: opensearch +status: + loadBalancer: {} + diff --git a/functions/k8s-confs/orborus-deployment.yaml b/functions/k8s-confs/orborus-deployment.yaml new file mode 100644 index 00000000..3725aa6d --- /dev/null +++ b/functions/k8s-confs/orborus-deployment.yaml @@ -0,0 +1,59 @@ +apiVersion: apps/v1 +kind: Deployment +metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.service: orborus + name: orborus +spec: + replicas: 1 + selector: + matchLabels: + io.kompose.service: orborus + strategy: {} + template: + metadata: + annotations: + kompose.cmd: kompose convert -f docker-compose.yml + kompose.version: 1.26.0 (40646f47) + creationTimestamp: null + labels: + io.kompose.network/shuffle: "true" + io.kompose.service: orborus + spec: + volumes: + - name: docker-socket + hostPath: + path: /var/run/docker.sock + containers: + - env: + - name: BASE_URL + value: "http://192.168.49.2:30009" + - name: DOCKER_API_VERSION + value: "1.40" + - name: ENVIRONMENT_NAME + value: Shuffle + - name: ORG_ID + value: Shuffle + - name: SHUFFLE_APP_SDK_VERSION + value: nightly + - name: SHUFFLE_SCALE_REPLICAS + value: "5" + #- name: SHUFFLE_SWARM_CONFIG + #value: run + - name: SHUFFLE_WORKER_VERSION + value: nightly + + image: orborus:v2 + imagePullPolicy: Never + name: shuffle-orborus + resources: {} + volumeMounts: + - name: docker-socket + mountPath: /var/run/docker.sock + hostname: shuffle-orborus + restartPolicy: Always +status: {} diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 21e42e9e..470d94db 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -36,11 +36,21 @@ import ( //"github.com/docker/docker/api/types/filters" dockerclient "github.com/docker/docker/client" - uuid "github.com/satori/go.uuid" + // uuid "github.com/satori/go.uuid" //"github.com/mackerelio/go-osstat/disk" "github.com/mackerelio/go-osstat/memory" "github.com/shirou/gopsutil/cpu" + + //k8s deps + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/tools/clientcmd" + "k8s.io/client-go/util/homedir" + "k8s.io/client-go/rest" + "path/filepath" + + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) // Starts jobs in bulk, so this could be increased @@ -559,6 +569,15 @@ func deployServiceWorkers(image string) { // Deploys the internal worker whenever something happens // https://docs.docker.com/engine/api/sdk/examples/ + +func buildEnvVars(envMap map[string]string) []corev1.EnvVar { + var envVars []corev1.EnvVar + for key, value := range envMap { + envVars = append(envVars, corev1.EnvVar{Name: key, Value: value}) + } + return envVars +} + func deployWorker(image string, identifier string, env []string, executionRequest shuffle.ExecutionRequest) error { // Binds is the actual "-v" volume. // Max 20% CPU every second @@ -566,157 +585,200 @@ func deployWorker(image string, identifier string, env []string, executionReques //CPUQuota: 25000, //CPUPeriod: 100000, //CPUShares: 256, - hostConfig := &container.HostConfig{ - LogConfig: container.LogConfig{ - Type: "json-file", - Config: map[string]string{ - "max-size": "10m", + // fmt.Printf("IMAGE NAME: %s",image) + // fmt.Printf("EEXECUTION REQUEST IN ORBORUS: %+v",executionRequest) + fmt.Printf("ENV: %+v", env) + image = "shuffle-worker:v1" + + 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 := getKubernetesClient() + if err != nil { + fmt.Println("[ERROR]Error getting kubernetes client:", err) + os.Exit(1) + } + + pod := &corev1.Pod{ + ObjectMeta: metav1.ObjectMeta{ + Name: identifier, + }, + Spec: corev1.PodSpec{ + Containers: []corev1.Container{ + { + Name: identifier, + Image: image, + Env: buildEnvVars(envMap), + }, + }, }, - }, - Resources: container.Resources{}, - } - - if len(os.Getenv("DOCKER_HOST")) == 0 { - if runtime.GOOS == "windows" { - hostConfig.Binds = []string{`\\.\pipe\docker_engine:\\.\pipe\docker_engine`} - } else { - hostConfig.Binds = []string{"/var/run/docker.sock:/var/run/docker.sock:rw"} } - } - - hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId)) - - if strings.ToLower(cleanupEnv) == "true" { - hostConfig.AutoRemove = true - } - - config := &container.Config{ - Image: image, - Env: env, - } - - //var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG") - parsedUuid := uuid.NewV4() - if swarmConfig == "run" || swarmConfig == "swarm" { - // FIXME: Should we handle replies properly? - // In certain cases, a workflow may e.g. be aborted already. If it's aborted, that returns - // a 401 from the worker, which returns an error here - go sendWorkerRequest(executionRequest) - //sendWorkerRequest(executionRequest) - - //err := sendWorkerRequest(executionRequest) - //if err != nil { - // log.Printf("[ERROR] Failed worker request for %s: %s", executionRequest.ExecutionId, err) - - // if strings.Contains(fmt.Sprintf("%s", err), "connection refused") || strings.Contains(fmt.Sprintf("%s", err), "EOF") { - // workerImage := fmt.Sprintf("%s/%s/shuffle-worker:%s", baseimageregistry, baseimagename, workerVersion) - // deployServiceWorkers(workerImage) - - // time.Sleep(time.Duration(10) * time.Second) - // err = sendWorkerRequest(executionRequest) - // } - //} - - //if err == nil { - // // FIXME: Readd this? Removed for rerun reasons - // // executionIds = append(executionIds, executionRequest.ExecutionId) - //} - //}() - - return nil - } - - //log.Printf("[INFO] Identifier: %s", identifier) - cont, err := dockercli.ContainerCreate( - context.Background(), - config, - hostConfig, - nil, - nil, - identifier, - ) - - if err != nil { - if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") { - identifier = fmt.Sprintf("%s-%s", identifier, parsedUuid) - //log.Printf("[INFO] 2 - Identifier: %s", identifier) - cont, err = dockercli.ContainerCreate( - context.Background(), - config, - hostConfig, - nil, - nil, - identifier, - ) - - if err != nil { - log.Printf("[ERROR] Container create error(2): %s", err) - return err - } - } else { - log.Printf("[ERROR] Container create error: %s", err) - return err - } - } - - containerStartOptions := types.ContainerStartOptions{} - err = dockercli.ContainerStart(context.Background(), cont.ID, containerStartOptions) - if err != nil { - // Trying to recreate and start WITHOUT network if it's possible. No extended checks. Old execution system (<0.9.30) - if strings.Contains(fmt.Sprintf("%s", err), "cannot join network") || strings.Contains(fmt.Sprintf("%s", err), "No such container") { - hostConfig.NetworkMode = "" - //container.NetworkMode(fmt.Sprintf("container:%s", containerId)) - cont, err = dockercli.ContainerCreate( - context.Background(), - config, - hostConfig, - nil, - nil, - identifier+"-2", - ) - if err != nil { - log.Printf("[ERROR] Failed to CREATE container (2): %s", err) - } - - err = dockercli.ContainerStart(context.Background(), cont.ID, containerStartOptions) - if err != nil { - log.Printf("[ERROR] Failed to start container (2): %s", err) - } - } else { - log.Printf("[ERROR] Failed initial container start. Quitting as this is NOT a simple network issue. Err: %s", err) - } - + + createdPod, err := clientset.CoreV1().Pods("shuffle").Create(context.Background(), pod, metav1.CreateOptions{}) if err != nil { - log.Printf("[ERROR] Failed to start worker container in environment %s: %s", environment, err) - return err - } else { - log.Printf("[INFO] Worker Container %s was created under environment %s for execution %s: docker logs %s", cont.ID, environment, executionRequest.ExecutionId, cont.ID) + fmt.Fprintf(os.Stderr, "Error creating pod: %v\n", err) + os.Exit(1) } + + fmt.Printf("Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace) + ///////////////////////////////////////////// - //stats, err := cli.ContainerInspect(context.Background(), containerName) - //if err != nil { - // log.Printf("Failed checking worker %s", containerName) - // return - //} + // hostConfig := &container.HostConfig{ + // LogConfig: container.LogConfig{ + // Type: "json-file", + // Config: map[string]string{ + // "max-size": "10m", + // }, + // }, + // Resources: container.Resources{}, + // } - //containerStatus := stats.ContainerJSONBase.State.Status - //if containerStatus != "running" { - // log.Printf("Status of %s is %s. Should be running. Will reset", containerName, containerStatus) - // err = stopWorker(containerName) - // if err != nil { - // log.Printf("Failed stopping worker %s", execution.ExecutionId) - // return - // } + // if len(os.Getenv("DOCKER_HOST")) == 0 { + // if runtime.GOOS == "windows" { + // hostConfig.Binds = []string{`\\.\pipe\docker_engine:\\.\pipe\docker_engine`} + // } else { + // hostConfig.Binds = []string{"/var/run/docker.sock:/var/run/docker.sock:rw"} + // } + // } - // err = deployWorke(cli, workerImage, containerName, env) - // if err != nil { - // log.Printf("Failed executing worker %s in state %s", execution.ExecutionId, containerStatus) - // return - // } - //} - } else { - log.Printf("[INFO] Worker Container %s was created under environment %s: docker logs %s", cont.ID, environment, cont.ID) - } + // hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:%s", containerId)) + + // if strings.ToLower(cleanupEnv) == "true" { + // hostConfig.AutoRemove = true + // } + + // config := &container.Config{ + // Image: image, + // Env: env, + // } + + // //var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG") + // parsedUuid := uuid.NewV4() + // if swarmConfig == "run" || swarmConfig == "swarm" { + // // FIXME: Should we handle replies properly? + // // In certain cases, a workflow may e.g. be aborted already. If it's aborted, that returns + // // a 401 from the worker, which returns an error here + // go sendWorkerRequest(executionRequest) + // //sendWorkerRequest(executionRequest) + + // //err := sendWorkerRequest(executionRequest) + // //if err != nil { + // // log.Printf("[ERROR] Failed worker request for %s: %s", executionRequest.ExecutionId, err) + + // // if strings.Contains(fmt.Sprintf("%s", err), "connection refused") || strings.Contains(fmt.Sprintf("%s", err), "EOF") { + // // workerImage := fmt.Sprintf("%s/%s/shuffle-worker:%s", baseimageregistry, baseimagename, workerVersion) + // // deployServiceWorkers(workerImage) + + // // time.Sleep(time.Duration(10) * time.Second) + // // err = sendWorkerRequest(executionRequest) + // // } + // //} + + // //if err == nil { + // // // FIXME: Readd this? Removed for rerun reasons + // // // executionIds = append(executionIds, executionRequest.ExecutionId) + // //} + // //}() + + // return nil + // } + + // //log.Printf("[INFO] Identifier: %s", identifier) + // cont, err := dockercli.ContainerCreate( + // context.Background(), + // config, + // hostConfig, + // nil, + // nil, + // identifier, + // ) + + // if err != nil { + // if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") { + // identifier = fmt.Sprintf("%s-%s", identifier, parsedUuid) + // //log.Printf("[INFO] 2 - Identifier: %s", identifier) + // cont, err = dockercli.ContainerCreate( + // context.Background(), + // config, + // hostConfig, + // nil, + // nil, + // identifier, + // ) + + // if err != nil { + // log.Printf("[ERROR] Container create error(2): %s", err) + // return err + // } + // } else { + // log.Printf("[ERROR] Container create error: %s", err) + // return err + // } + // } + + // containerStartOptions := types.ContainerStartOptions{} + // err = dockercli.ContainerStart(context.Background(), cont.ID, containerStartOptions) + // if err != nil { + // // Trying to recreate and start WITHOUT network if it's possible. No extended checks. Old execution system (<0.9.30) + // if strings.Contains(fmt.Sprintf("%s", err), "cannot join network") || strings.Contains(fmt.Sprintf("%s", err), "No such container") { + // hostConfig.NetworkMode = "" + // //container.NetworkMode(fmt.Sprintf("container:%s", containerId)) + // cont, err = dockercli.ContainerCreate( + // context.Background(), + // config, + // hostConfig, + // nil, + // nil, + // identifier+"-2", + // ) + // if err != nil { + // log.Printf("[ERROR] Failed to CREATE container (2): %s", err) + // } + + // err = dockercli.ContainerStart(context.Background(), cont.ID, containerStartOptions) + // if err != nil { + // log.Printf("[ERROR] Failed to start container (2): %s", err) + // } + // } else { + // log.Printf("[ERROR] Failed initial container start. Quitting as this is NOT a simple network issue. Err: %s", err) + // } + + // if err != nil { + // log.Printf("[ERROR] Failed to start worker container in environment %s: %s", environment, err) + // return err + // } else { + // log.Printf("[INFO] Worker Container %s was created under environment %s for execution %s: docker logs %s", cont.ID, environment, executionRequest.ExecutionId, cont.ID) + // } + + // //stats, err := cli.ContainerInspect(context.Background(), containerName) + // //if err != nil { + // // log.Printf("Failed checking worker %s", containerName) + // // return + // //} + + // //containerStatus := stats.ContainerJSONBase.State.Status + // //if containerStatus != "running" { + // // log.Printf("Status of %s is %s. Should be running. Will reset", containerName, containerStatus) + // // err = stopWorker(containerName) + // // if err != nil { + // // log.Printf("Failed stopping worker %s", execution.ExecutionId) + // // return + // // } + + // // err = deployWorke(cli, workerImage, containerName, env) + // // if err != nil { + // // log.Printf("Failed executing worker %s in state %s", execution.ExecutionId, containerStatus) + // // return + // // } + // //} + // } else { + // log.Printf("[INFO] Worker Container %s was created under environment %s: docker logs %s", cont.ID, environment, cont.ID) + // } return nil } @@ -930,8 +992,47 @@ func getOrborusStats() shuffle.OrborusStats { return newStats } +func isRunningInCluster() bool { + _, existsHost := os.LookupEnv("KUBERNETES_SERVICE_HOST") + _, existsPort := os.LookupEnv("KUBERNETES_SERVICE_PORT") + return existsHost && existsPort +} + +func getKubernetesClient() (*kubernetes.Clientset, error) { + if isRunningInCluster() { + config, err := rest.InClusterConfig() + if err != nil { + return nil, err + } + clientset, err := kubernetes.NewForConfig(config) + if err != nil { + return nil, err + } + return clientset, nil + } else { + home := homedir.HomeDir() + kubeconfigPath := filepath.Join(home, ".kube", "config") + config, err := clientcmd.BuildConfigFromFlags("", kubeconfigPath) + if err != nil { + return nil, err + } + clientset, err := kubernetes.NewForConfig(config) + if err != nil { + return nil, err + } + return clientset, nil + } +} + // Initial loop etc func main() { + + if isRunningInCluster(){ + log.Printf("[INFO] Running inside k8s cluster") + } + + ///////////////////////// + startupDelay := os.Getenv("SHUFFLE_ORBORUS_STARTUP_DELAY") if len(startupDelay) > 0 { log.Printf("[DEBUG] Setting startup delay to %#v", startupDelay) @@ -946,7 +1047,7 @@ func main() { log.Println("[INFO] Setting up execution environment") - //FIXME + // //FIXME if baseUrl == "" { baseUrl = "https://shuffler.io" //baseUrl = "http://localhost:5001" @@ -1015,7 +1116,8 @@ func main() { ctx := context.Background() // Run by default from now - zombiecheck(ctx, workerTimeout) + //commenting for now as its stoppoing minikube + // zombiecheck(ctx, workerTimeout) log.Printf("[INFO] Running towards %s (BASE_URL) with environment name %s", baseUrl, environment) diff --git a/functions/onprem/worker/Dockerfile b/functions/onprem/worker/Dockerfile index 3b89d290..554f2357 100755 --- a/functions/onprem/worker/Dockerfile +++ b/functions/onprem/worker/Dockerfile @@ -32,6 +32,8 @@ FROM alpine:3.15.0 ENV SHUFFLE_BASE_IMAGE_REGISTRY=docker.io ENV SHUFFLE_BASE_IMAGE_NAME=frikky/shuffle ENV SHUFFLE_BASE_IMAGE_TAG_SUFFIX=0.8.70 +ENV SHUFFLE_OPENSEARCH_URL=https://192.168.49.2:31001 +ENV SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY=true RUN apk add --no-cache bash tzdata COPY --from=builder /app/ /