diff --git a/.github/workflows/helm-release.yml b/.github/workflows/helm-release.yml index 4c3c6b6f..69b17917 100644 --- a/.github/workflows/helm-release.yml +++ b/.github/workflows/helm-release.yml @@ -26,9 +26,9 @@ jobs: - name: Install apt dependencies run: | - curl https://baltocdn.com/helm/signing.asc | gpg --dearmor | sudo tee /usr/share/keyrings/helm.gpg > /dev/null sudo apt-get install apt-transport-https -y --no-install-recommends - echo "deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/helm.gpg] https://baltocdn.com/helm/stable/debian/ all main" | sudo tee /etc/apt/sources.list.d/helm-stable-debian.list + curl -fsSL https://packages.buildkite.com/helm-linux/helm-debian/gpgkey | gpg --dearmor | sudo tee /usr/share/keyrings/helm.gpg > /dev/null + echo "deb [signed-by=/usr/share/keyrings/helm.gpg] https://packages.buildkite.com/helm-linux/helm-debian/any/ any main" | sudo tee /etc/apt/sources.list.d/helm-stable-debian.list sudo apt-get update sudo apt-get install helm -y --no-install-recommends diff --git a/.github/workflows/helm-test.yml b/.github/workflows/helm-test.yml index fb0e34aa..49c20240 100644 --- a/.github/workflows/helm-test.yml +++ b/.github/workflows/helm-test.yml @@ -18,9 +18,9 @@ jobs: - name: Install apt dependencies run: | - curl https://baltocdn.com/helm/signing.asc | gpg --dearmor | sudo tee /usr/share/keyrings/helm.gpg > /dev/null sudo apt-get install apt-transport-https -y --no-install-recommends - echo "deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/helm.gpg] https://baltocdn.com/helm/stable/debian/ all main" | sudo tee /etc/apt/sources.list.d/helm-stable-debian.list + curl -fsSL https://packages.buildkite.com/helm-linux/helm-debian/gpgkey | gpg --dearmor | sudo tee /usr/share/keyrings/helm.gpg > /dev/null + echo "deb [signed-by=/usr/share/keyrings/helm.gpg] https://packages.buildkite.com/helm-linux/helm-debian/any/ any main" | sudo tee /etc/apt/sources.list.d/helm-stable-debian.list sudo apt-get update sudo apt-get install helm -y --no-install-recommends diff --git a/backend/go-app/main.go b/backend/go-app/main.go index 17817cd5..f1d53d63 100755 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -69,11 +69,12 @@ var debug = false //var syncUrl = "http://localhost:5002" type retStruct struct { - Success bool `json:"success"` - SyncFeatures shuffle.SyncFeatures `json:"sync_features"` - SessionKey string `json:"session_key"` - IntervalSeconds int64 `json:"interval_seconds"` - Reason string `json:"reason"` + Success bool `json:"success"` + SyncFeatures shuffle.SyncFeatures `json:"sync_features"` + SessionKey string `json:"session_key"` + IntervalSeconds int64 `json:"interval_seconds"` + Reason string `json:"reason"` + Subscriptions []shuffle.PaymentSubscription `json:"subscriptions"` } type Contact struct { @@ -3870,10 +3871,11 @@ func handleCloudJob(job shuffle.CloudSyncJob) error { // Handles jobs from remote (cloud) func remoteOrgJobController(org shuffle.Org, body []byte) error { type retStruct struct { - Success bool `json:"success"` - Reason string `json:"reason"` - Jobs []shuffle.CloudSyncJob `json:"jobs"` + Success bool `json:"success"` + Reason string `json:"reason"` + Jobs []shuffle.CloudSyncJob `json:"jobs"` SyncFeatures shuffle.SyncFeatures `json:"sync_features"` + Subscriptions []shuffle.PaymentSubscription `json:"subscriptions"` } responseData := retStruct{} @@ -4711,6 +4713,7 @@ func handleStopCloudSync(syncUrl string, org shuffle.Org) (*shuffle.Org, error) org.CloudSync = false org.SyncFeatures = shuffle.SyncFeatures{} org.SyncConfig = shuffle.SyncConfig{} + org.Subscriptions = []shuffle.PaymentSubscription{} err = shuffle.SetOrg(ctx, org, org.Id) if err != nil { @@ -4926,8 +4929,6 @@ func handleCloudSetup(resp http.ResponseWriter, request *http.Request) { return } - log.Printf("[DEBUG] Respbody from sync: %s", string(respBody)) - responseData := retStruct{} err = json.Unmarshal(respBody, &responseData) if err != nil { diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 3cf355ab..090d524e 100755 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -1847,13 +1847,12 @@ func getWorkflowApps(resp http.ResponseWriter, request *http.Request) { workflowapps, err := shuffle.GetAllWorkflowApps(ctx, 1000, 0) if err != nil { log.Printf("{WARNING] Failed getting apps (getworkflowapps): %s", err) - resp.WriteHeader(401) + resp.WriteHeader(400) resp.Write([]byte(`{"success": false}`)) return } newapps := workflowapps - if len(user.PrivateApps) > 0 { found := false for _, item := range user.PrivateApps { diff --git a/docker-compose.yml b/docker-compose.yml index c1c27efb..8ab7e3b9 100755 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -91,6 +91,27 @@ services: - shuffle restart: unless-stopped + # OPTIONAL: advanced monitoring with cAdvisor + # If you have a seperate orborus config YAML on another server, + # Please add it there as well if you want to monitor that server too. + # FYI: both servers will show up on different dashboards. + # cadvisor: + # image: gcr.io/cadvisor/cadvisor:latest + # volumes: + # - /:/rootfs:ro + # - /var/run:/var/run:ro + # - /sys:/sys:ro + # - /var/lib/docker/:/var/lib/docker:ro + # - /dev/disk/:/dev/disk:ro + # - /var/run/docker.sock:/var/run/docker.sock:ro + # ports: + # - "8080:8080" + # privileged: true + # devices: + # - /dev/kmsg:/dev/kmsg + # networks: + # - shuffle + #memcached: # image: memcached:latest # container_name: shuffle-cache diff --git a/frontend/src/codeeditor-index.css b/frontend/src/codeeditor-index.css index c87341f5..9ccc0acd 100644 --- a/frontend/src/codeeditor-index.css +++ b/frontend/src/codeeditor-index.css @@ -1,9 +1,7 @@ -@import url('https://fonts.googleapis.com/css?family=Nunito+Sans'); - body { margin: 0; padding: 0; - font-family: "Nunito Sans", sans-serif; + font-family: "Segoe UI", Roboto, "Noto Sans", "Liberation Sans", Arial, "Helvetica Neue", -apple-system, BlinkMacSystemFont, "Nunito Sans", sans-serif; -webkit-font-smoothing: antialiased; -moz-osx-font-smoothing: grayscale; } diff --git a/functions/kubernetes/README.md b/functions/kubernetes/README.md index a54884e8..c2d2c3d1 100644 --- a/functions/kubernetes/README.md +++ b/functions/kubernetes/README.md @@ -1,68 +1,12 @@ -## How to deploy Shuffle on Kubernetes? +# Shuffle in Kubernetes -### Prerequisites: -- Clone the https://github.com/shuffle/shuffle repository using Git then, navigate to the functions/kubernetes directory, which contains all the necessary Kubernetes configuration files for deployment. -- [Running a Kubernetes cluster](https://kubernetes.io/docs/setup/). You can do that with either minikube or run the cluster locally. -- Ensure you have a local Docker registry set up to store and manage Docker images for applications built with Shuffle. While the registry is crucial for handling custom-built apps, you’ll still be able to run workflows without it. To setup a docker registry, if you have docker installed on one of your node run following commands. - - - ``` - chmod +x generate_certs.sh - ./setup_registry.sh - ``` - - > This will give you a NODE_IP which is you're local IP if you're not sure about what to use. - - > **Make sure that port 5000 is not exposed to the internet!** - -- 8 GB RAM and 4 CPUs are recommended as **minimum configs** for running Shuffle on Kubernetes. K8s is a resource-intensive application, and you may experience performance issues if you run it on a machine with fewer resources. - -- If you've used the above commands to set up a registry, you'll need to skip an SSL verification for your registry. If you're using Containerd as a runtime - add the following lines in /etc/containerd/config.toml - ``` - [plugins."io.containerd.grpc.v1.cri".registry.mirrors.""] - endpoint = ["https://"] - - [plugins."io.containerd.grpc.v1.cri".registry.configs."".tls] - insecure_skip_verify = true - ``` - -### Instructions -Step 1: Create a namespace called shuffle in a cluster by running ```kubectl create ns shuffle```. - -Step 2: Open the ```all-in-one.yaml``` file and review the configuration values. Change the value of REGISTRY_URL with ':5000' where the registry is at. Adjust other variables as per your deployment requirements; otherwise, the application will deploy using the default settings provided within the file. Then apply the configmap and deploy with ```kubectl apply -f all-in-one.yaml -n shuffle``` - -Step 3: Now, open ```https://:30008``` or ```http://:30007```. You should be seeing a signup page. NODE_IP should be where the frontend is deployed. - -### Dev Mode - -1. Run backend and orborus with the environment variable `IS_KUBERNETES=true`: - -```bash -export IS_KUBERNETES=true -``` - -2. Turn on the k8s engine with minikube: - -```bash -minikube start -``` - -3. To use the worker scale feature, build the image with the following command: - -```bash -$NAME=shuffle-worker-scale -$VERSION=1.2.0 - -minikube build . -t shuffle/shuffle:$NAME -t shuffle/shuffle:$NAME_$VERSION -t docker.pkg.github.com/shuffle/shuffle/$NAME:$VERSION -t ghcr.io/shuffle/$NAME:$VERSION -t ghcr.io/shuffle/$NAME:nightly -t ghcr.io/shuffle/$NAME:$VERSION -t ghcr.io/shuffle/$NAME:nightly -t ghcr.io/shuffle/$NAME:latest -``` - -4. To run executions, Make sure to do the following: - -```bash -kubectl create role pod-creator --namespace=default --verb=create --resource=pods -kubectl create rolebinding pod-creator-binding --namespace=default --role=pod-creator --serviceaccount=default:default -``` +1. Make sure you have a Kubernetes cluster available. MiniKube works for testing. +2. Install `helm install shuffle oci://ghcr.io/shuffle/charts/shuffle --namespace shuffle --create-namespace` +3. Tweak the configuration files if needed! This is not meant to be a one-size-fits-all +More details in the [kubernetes/Charts/Shuffle folder.](https://github.com/Shuffle/Shuffle/tree/main/functions/kubernetes/charts/shuffle#usage) +## Architecture +Here is the default architecture it follows, with the "Frontend" being the exposed container you interact with. +image diff --git a/functions/kubernetes/all-in-one.yaml b/functions/kubernetes/all-in-one.yaml deleted file mode 100644 index d9bec68d..00000000 --- a/functions/kubernetes/all-in-one.yaml +++ /dev/null @@ -1,894 +0,0 @@ -apiVersion: v1 -kind: Namespace -metadata: - name: shuffle - ---- -## ONLY for minikube -# apiVersion: storage.k8s.io/v1 -# kind: StorageClass -# metadata: -# name: standard-rwo -# provisioner: k8s.io/minikube-hostpath -# reclaimPolicy: Delete -# volumeBindingMode: Immediate - -# --- - -apiVersion: storage.k8s.io/v1 -kind: StorageClass -metadata: - name: shuffle-data - namespace: shuffle -provisioner: kubernetes.io/no-provisioner -volumeBindingMode: WaitForFirstConsumer - ---- - -apiVersion: v1 -metadata: - namespace: shuffle - creationTimestamp: null - labels: - io.kompose.service: backend-env - name: env -data: - BACKEND_HOSTNAME: shuffle-backend - BACKEND_PORT: "5001" - BASE_URL: http://shuffle-backend:5001 - DATASTORE_EMULATOR_HOST: shuffle-database:8000 - DB_LOCATION: /mnt/shuffle-data/open-search - DOCKER_API_VERSION: "1.40" - ENVIRONMENT_NAME: Shuffle - FRONTEND_PORT: "3001" - FRONTEND_PORT_HTTPS: "3443" - HTTP_PROXY: "" - HTTPS_PROXY: "" - ORBORUS_CONTAINER_NAME: "\t\t\t\t" - ORG_ID: Shuffle - OUTER_HOSTNAME: shuffle-backend - SHUFFLE_APP_DOWNLOAD_LOCATION: https://github.com/shuffle/python-apps - SHUFFLE_APP_FORCE_UPDATE: "false" - SHUFFLE_APP_HOTLOAD_FOLDER: /shuffle-apps - SHUFFLE_APP_HOTLOAD_LOCATION: ./shuffle-apps - SHUFFLE_BASE_IMAGE_NAME: shuffle - SHUFFLE_BASE_IMAGE_REGISTRY: ghcr.io - SHUFFLE_BASE_IMAGE_TAG_SUFFIX: -1.0.0 - SHUFFLE_CHAT_DISABLED: "false" - SHUFFLE_CONTAINER_AUTO_CLEANUP: "false" - SHUFFLE_DEFAULT_APIKEY: "" - SHUFFLE_DEFAULT_PASSWORD: "" - SHUFFLE_DEFAULT_USERNAME: "" - SHUFFLE_DOWNLOAD_AUTH_BRANCH: "" - SHUFFLE_DOWNLOAD_AUTH_PASSWORD: "" - SHUFFLE_DOWNLOAD_AUTH_USERNAME: "" - SHUFFLE_DOWNLOAD_WORKFLOW_BRANCH: "" - SHUFFLE_DOWNLOAD_WORKFLOW_LOCATION: "" - SHUFFLE_DOWNLOAD_WORKFLOW_PASSWORD: "" - SHUFFLE_DOWNLOAD_WORKFLOW_USERNAME: "" - SHUFFLE_ELASTIC: "true" - SHUFFLE_ENCRYPTION_MODIFIER: "" - SHUFFLE_FILE_LOCATION: /shuffle-files - SHUFFLE_LOGS_DISABLED: "false" - SHUFFLE_OPENSEARCH_APIKEY: "" - SHUFFLE_OPENSEARCH_CERTIFICATE_FILE: "" - SHUFFLE_OPENSEARCH_CLOUDID: "" - KUBERNETES_NAMESPACE: shuffle - SHUFFLE_OPENSEARCH_INDEX_PREFIX: "" - SHUFFLE_OPENSEARCH_PASSWORD: admin - SHUFFLE_OPENSEARCH_PROXY: "" - SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY: "true" - SHUFFLE_OPENSEARCH_URL: https://opensearch:9200 - SHUFFLE_MEMCACHED: shuffle-memcached:11211 - SHUFFLE_OPENSEARCH_USERNAME: admin - SHUFFLE_ORBORUS_STARTUP_DELAY: "\t\t" - SHUFFLE_PASS_APP_PROXY: "FALSE" - SHUFFLE_PASS_WORKER_PROXY: "TRUE" - SHUFFLE_RERUN_SCHEDULE: "300" - SSO_REDIRECT_URL: "" - TZ: "Europe/Amsterdam \t\t\t\t\t" - IS_KUBERNETES: "true" - REGISTRY_URL: "docker-registry:5000" - REGISTRY_AUTH: "false" - SHUFFLE_KUBERNETES_WORKER: "ghcr.io/shuffle/shuffle-worker:latest" -kind: ConfigMap - ---- -apiVersion: rbac.authorization.k8s.io/v1 -kind: Role -metadata: - namespace: shuffle - name: pod-manager -rules: -- apiGroups: [""] - resources: ["pods", "services", "deployments"] - verbs: ["get", "list", "create", "update", "delete"] -- apiGroups: ["batch"] - resources: ["jobs"] - verbs: ["create", "get", "list", "watch", "delete"] -- apiGroups: ["rbac.authorization.k8s.io"] - resources: ["rolebindings", "roles"] - verbs: ["get", "list", "create"] -- apiGroups: ["apps"] - resources: ["deployments", "pods", "services"] - verbs: ["create", "get", "list", "update", "delete"] - ---- - -apiVersion: rbac.authorization.k8s.io/v1 -kind: RoleBinding -metadata: - name: pod-manager-binding - namespace: shuffle -subjects: -- kind: ServiceAccount - name: default - namespace: shuffle -roleRef: - kind: Role - name: pod-manager - apiGroup: rbac.authorization.k8s.io - ---- - -# apiVersion: v1 -# kind: PersistentVolume -# metadata: -# name: shuffle-os-pv -# namespace: shuffle -# 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: standard-rwo # Set the desired storage class -# hostPath: -# path: /mnt/shuffle-data/open-search - -# --- - -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - namespace: shuffle - creationTimestamp: null - labels: - io.kompose.service: opensearch-claim0 - name: opensearch-claim0 -spec: - accessModes: - - ReadWriteOnce - storageClassName: standard-rwo - resources: - requests: - storage: 500Mi -status: {} - ---- -apiVersion: apps/v1 -kind: Deployment -metadata: - namespace: shuffle - 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: - # securityContext: - # runAsUser: 1000 # UID - # fsGroup: 1000 # GID - # nodeSelector: - # node: worker1 - initContainers: - - name: volume-permissions - image: busybox - command: ["sh", "-c", "chown -R 1000:1000 /usr/share/opensearch/data"] - volumeMounts: - - name: opensearch-claim0 - mountPath: /usr/share/opensearch/data - 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: - namespace: shuffle - 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: - ports: - - name: "9200" - port: 9200 - targetPort: 9200 - selector: - io.kompose.service: opensearch -status: - loadBalancer: {} - ---- - -# apiVersion: v1 -# kind: PersistentVolume -# metadata: -# namespace: shuffle -# name: shuffle-apps-pv -# spec: -# capacity: -# storage: 5Gi -# accessModes: -# - ReadWriteOnce -# persistentVolumeReclaimPolicy: Retain -# storageClassName: shuffle-data -# hostPath: -# path: /mnt/shuffle-data/backend - -# --- -# apiVersion: v1 -# kind: PersistentVolume -# metadata: -# namespace: shuffle -# name: shuffle-files-pv -# spec: -# capacity: -# storage: 5Gi -# accessModes: -# - ReadWriteOnce -# persistentVolumeReclaimPolicy: Retain -# storageClassName: shuffle-data -# hostPath: -# path: /mnt/shuffle-data/backend - -# --- - -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - namespace: shuffle - creationTimestamp: null - labels: - io.kompose.service: backend-files-claim - name: backend-files-claim -spec: - accessModes: - - ReadWriteOnce - storageClassName: standard-rwo - resources: - requests: - storage: 5Gi -# status: {} - ---- - -apiVersion: v1 -kind: PersistentVolumeClaim -metadata: - namespace: shuffle - creationTimestamp: null - labels: - io.kompose.service: backend-apps-claim - name: backend-apps-claim -spec: - accessModes: - - ReadWriteOnce - storageClassName: standard-rwo - resources: - requests: - storage: 5Gi -# status: {} - ---- - -apiVersion: apps/v1 -kind: Deployment -metadata: - name: shuffle-memcached - namespace: shuffle -spec: - replicas: 1 - selector: - matchLabels: - app: shuffle-memcached - template: - metadata: - labels: - app: shuffle-memcached - spec: - containers: - - name: shuffle-memcached - image: memcached:latest - ports: - - containerPort: 11211 - resources: {} - restartPolicy: Always - - ---- - -apiVersion: v1 -kind: Service -metadata: - namespace: shuffle - name: shuffle-memcached -spec: - ports: - - port: 11211 - targetPort: 11211 - selector: - app: shuffle-memcached - type: ClusterIP - ---- - -apiVersion: apps/v1 -kind: Deployment -metadata: - namespace: shuffle - 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 -# nodeSelector: -# node: master - 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_MEMCACHED - valueFrom: - configMapKeyRef: - key: SHUFFLE_MEMCACHED - 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 - - name: IS_KUBERNETES - valueFrom: - configMapKeyRef: - key: IS_KUBERNETES - name: env - - name: REGISTRY_URL - valueFrom: - configMapKeyRef: - key: REGISTRY_URL - name: env - - name: REGISTRY_AUTH - valueFrom: - configMapKeyRef: - key: REGISTRY_AUTH - name: env - image: ghcr.io/shuffle/shuffle-backend:latest - imagePullPolicy: Always - name: shuffle-backend - ports: - - containerPort: 5001 - resources: {} - volumeMounts: - - name: shuffle-apps - mountPath: /app/generated - - name: shuffle-files - mountPath: /shuffle-files - restartPolicy: Always -status: {} - ---- - -apiVersion: v1 -kind: Service -metadata: - namespace: shuffle - 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: - ports: - - name: "5001" - port: 5001 - targetPort: 5001 - selector: - io.kompose.service: backend -status: - loadBalancer: {} - ---- -apiVersion: apps/v1 -kind: Deployment -metadata: - namespace: shuffle - 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: - - name: shuffle-frontend - image: ghcr.io/shuffle/shuffle-frontend:latest - env: - - name: BACKEND_HOSTNAME - valueFrom: - configMapKeyRef: - key: BACKEND_HOSTNAME - name: env - ports: - - containerPort: 80 - - containerPort: 443 - resources: {} - hostname: shuffle-frontend - restartPolicy: Always -status: {} - ---- - -apiVersion: v1 -kind: Service -metadata: - namespace: shuffle - 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: {} - ---- - -apiVersion: apps/v1 -kind: Deployment -metadata: - namespace: shuffle - 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: - containers: - - env: - - name: BASE_URL - value: "http://shuffle-backend:5001" - - name: DOCKER_API_VERSION - value: "1.40" - - name: ENVIRONMENT_NAME - value: Shuffle - - name: ORG_ID - value: Shuffle - - name: SHUFFLE_APP_SDK_VERSION - value: latest - - name: SHUFFLE_SCALE_REPLICAS - value: "5" - - name: SHUFFLE_SWARM_CONFIG - value: run - - name: SHUFFLE_WORKER_VERSION - value: latest - - name: IS_KUBERNETES - valueFrom: - configMapKeyRef: - key: IS_KUBERNETES - name: env - - name: KUBERNETES_NAMESPACE - valueFrom: - configMapKeyRef: - key: KUBERNETES_NAMESPACE - name: env - - name: REGISTRY_URL - valueFrom: - configMapKeyRef: - key: REGISTRY_URL - name: env - - name: SHUFFLE_KUBERNETES_WORKER - valueFrom: - configMapKeyRef: - key: SHUFFLE_KUBERNETES_WORKER - name: env - - name: SHUFFLE_MEMCACHED - valueFrom: - configMapKeyRef: - key: SHUFFLE_MEMCACHED - name: env - image: ghcr.io/shuffle/shuffle-orborus:latest - #imagePullPolicy: Never - name: shuffle-orborus - resources: {} - hostname: shuffle-orborus - restartPolicy: Always -status: {} diff --git a/functions/kubernetes/charts/shuffle/Chart.yaml b/functions/kubernetes/charts/shuffle/Chart.yaml index b01df7da..e286e11f 100644 --- a/functions/kubernetes/charts/shuffle/Chart.yaml +++ b/functions/kubernetes/charts/shuffle/Chart.yaml @@ -2,13 +2,13 @@ apiVersion: v2 name: shuffle description: A Helm chart for deploying Shuffle on Kubernetes type: application -version: 0.0.0 # Set during publishing in GitHub actions +version: 2.1.0 # Set during publishing in GitHub actions appVersion: latest # Overwritten during publishing in GitHub actions dependencies: - name: common version: ^2.23.0 repository: oci://registry-1.docker.io/bitnamicharts - name: opensearch - version: ^1.3.0 - repository: oci://registry-1.docker.io/bitnamicharts + version: 2.0.10 + repository: https://charts.bitnami.com/bitnami condition: opensearch.enabled diff --git a/functions/kubernetes/charts/shuffle/values.yaml b/functions/kubernetes/charts/shuffle/values.yaml index 57bccb31..4f46e482 100644 --- a/functions/kubernetes/charts/shuffle/values.yaml +++ b/functions/kubernetes/charts/shuffle/values.yaml @@ -74,7 +74,7 @@ diagnosticMode: shuffle: ## @param shuffle.baseUrl The external base URL under which Shuffle is reachable. ## - baseUrl: "" + #baseUrl: "" ## ref: https://shuffler.io/docs/organizations ## This chart only supports single-tenant deployments at the moment @@ -525,6 +525,8 @@ backend: ## @param backend.openSearch.username The username that is used for authenticating with OpenSearch ## username: admin + + password: StrongShufflePassword321! ## @param backend.openSearch.certificateFile The path to a custom OpenSearch certificate file ## certificateFile: "" @@ -885,6 +887,16 @@ frontend: ## ref: https://kubernetes.io/docs/concepts/overview/working-with-objects/labels/ ## labels: {} + type: NodePort + ports: + - name: http + port: 80 + targetPort: 3001 + nodePort: 30080 + - name: https + port: 443 + targetPort: 3443 + nodePort: 30443 ## ServiceAccount configuration ## @@ -1187,7 +1199,17 @@ orborus: ## - name: FOO ## value: "bar" ## - extraEnvVars: [] + extraEnvVars: + - name: SHUFFLE_APP_SDK_TIMEOUT + value: "300" + - name: SHUFFLE_ORBORUS_EXCUTION_CONCURRENCY + value: "7" + - name: SHUFFLE_STATS_DISABLED + value: "true" + - name: KUBERNETES_NAMESPACE + value: "shuffle" + - name: SHUFFLE_BASE_IMAGE_NAME + value: "frikky/shuffle" ## @param orborus.extraEnvVarsCM Name of existing ConfigMap containing extra env vars for orborus containers ## extraEnvVarsCM: "" @@ -1832,7 +1854,7 @@ volumePermissions: ## image: registry: docker.io - repository: bitnami/os-shell + repository: bitnamilegacy/os-shell tag: 12-debian-12-r30 pullPolicy: IfNotPresent ## Optionally specify an array of imagePullSecrets. @@ -1888,16 +1910,25 @@ volumePermissions: ## opensearch: enabled: true + + sysctlImage: + enabled: false + + image: + registry: docker.io + repository: bitnamilegacy/opensearch + tag: "3.2.0" + master: replicaCount: 1 data: replicaCount: 1 coordinating: - replicaCount: 1 + replicaCount: 0 ingest: - replicaCount: 1 + replicaCount: 0 dashboards: - enabled: true + enabled: false ## @section Vault Parameters ## diff --git a/functions/kubernetes/generate_certs.sh b/functions/kubernetes/generate_certs.sh deleted file mode 100644 index cb316a72..00000000 --- a/functions/kubernetes/generate_certs.sh +++ /dev/null @@ -1,39 +0,0 @@ -#!/bin/bash -# Check if ifconfig is present and use it to get the default IP -if command -v ifconfig &> /dev/null; then - default_ip=$(ifconfig | grep 'inet ' | grep -v 127.0.0.1 | awk '{print $2}') -# Check if ip is present and use it if ifconfig is not available -elif command -v ip &> /dev/null; then - default_ip=$(ip addr show | grep -oP 'inet \K[\d.]+' | sed -n '2p') -# If both tools are not available, error out -else - echo "Error: Neither ifconfig nor ip command found in the machine. Exiting.." - exit 1 -fi - - -read -p "Enter your node IP to use for cert generation (default is $default_ip): " custom_ip -# Use localhost as the default value -node_ip=${custom_ip:-$default_ip} - -echo "Using node IP: $node_ip to generate SSL certs!" - -mkdir -p certs - -# Generate CA key -openssl req -newkey rsa:4096 -nodes -sha256 -keyout certs/reg.key -x509 -days 365 -out certs/reg.crt -subj "/CN=$node_ip" - -# generate a random string -random_string=$(openssl rand -hex 3) - -echo "Starting docker registry with name shuffle-local-registry-$random_string.." - -docker run -d -p 5000:5000 --restart=always --name "shuffle-local-registry-$random_string" \ - -v $(pwd)/certs:/certs \ - -e REGISTRY_HTTP_TLS_CERTIFICATE=/certs/reg.crt \ - -e REGISTRY_HTTP_TLS_KEY=/certs/reg.key \ - registry:2 - -echo "Set up certs and launched docker registry successfully!" - -echo "Please put $node_ip:5000 as the REGISTRY_URL in all-in-one.yaml file" \ No newline at end of file diff --git a/functions/kubernetes/orborus.yaml b/functions/kubernetes/orborus.yaml deleted file mode 100644 index 000d9fbd..00000000 --- a/functions/kubernetes/orborus.yaml +++ /dev/null @@ -1,80 +0,0 @@ ---- - -apiVersion: rbac.authorization.k8s.io/v1 -kind: Role -metadata: - namespace: default - name: pod-manager -rules: -- apiGroups: [""] - resources: ["pods"] - verbs: ["get", "list", "create", "update", "delete"] -- apiGroups: ["batch"] - resources: ["jobs"] - verbs: ["create", "get", "list", "watch", "delete"] - ---- - -apiVersion: rbac.authorization.k8s.io/v1 -kind: RoleBinding -metadata: - name: pod-manager-binding - namespace: default -subjects: -- kind: ServiceAccount - name: default - namespace: default -roleRef: - kind: Role - name: pod-manager - apiGroup: rbac.authorization.k8s.io - ---- - -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: - containers: - - env: - - name: BASE_URL - value: "https://shuffler.io" - - name: SHUFFLE_SCALE_REPLICAS - value: "7" - - name: IS_KUBERNETES - value: "true" - - name: ENVIRONMENT_NAME - value: "environment test" - - name: ORG - value: "9c938e5b-d812-40d9-92f0-93783f43ec0d" - - name: AUTH - value: "3663a270-bb3a-4678-a365-d879601a1a0c" - - image: ghcr.io/shuffle/shuffle-orborus:latest - #imagePullPolicy: Never - name: shuffle-orborus - resources: {} - hostname: shuffle-orborus - restartPolicy: Always diff --git a/functions/kubernetes/setup_registry.sh b/functions/kubernetes/setup_registry.sh deleted file mode 100755 index be4c266f..00000000 --- a/functions/kubernetes/setup_registry.sh +++ /dev/null @@ -1,36 +0,0 @@ -#!/bin/bash -# Check if ifconfig is present and use it to get the default IP -if command -v ifconfig &> /dev/null; then - default_ip=$(ifconfig | grep 'inet ' | grep -v 127.0.0.1 | awk '{print $2}') -# Check if ip is present and use it if ifconfig is not available -elif command -v ip &> /dev/null; then - default_ip=$(ip addr show | grep -oP 'inet \K[\d.]+' | sed -n '2p') -# If both tools are not available, error out -else - echo "Error: Neither ifconfig nor ip command found in the machine. Exiting.." - exit 1 -fi - - -read -p "Enter your node IP to use for cert generation (default is $default_ip): " custom_ip -# Use localhost as the default value -node_ip=${custom_ip:-$default_ip} - -echo "Using node IP: $node_ip to generate SSL certs!" - -mkdir -p certs - -# Generate CA key -openssl req -newkey rsa:4096 -nodes -sha256 -keyout certs/reg.key -x509 -days 365 -out certs/reg.crt -subj "/CN=$node_ip" - -echo "Generated certs/reg.key and certs/reg.crt!" - -echo -e "Please run:\n" -echo "docker run -d -p 5000:5000 --restart=always --name "shuffle-local-registry" \ - -v $(pwd)/certs:/certs \ - -e REGISTRY_HTTP_TLS_CERTIFICATE=/certs/reg.crt \ - -e REGISTRY_HTTP_TLS_KEY=/certs/reg.key \ - registry:2" - -echo -e "\nnow to start the reigstry :)!" - diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod index ed7e5add..b6ccae82 100644 --- a/functions/onprem/orborus/go.mod +++ b/functions/onprem/orborus/go.mod @@ -10,7 +10,7 @@ require ( github.com/docker/docker v28.3.3+incompatible github.com/docker/go-connections v0.5.0 github.com/satori/go.uuid v1.2.0 - github.com/shuffle/shuffle-shared v0.9.15 + github.com/shuffle/shuffle-shared v0.9.27 k8s.io/api v0.33.1 k8s.io/apimachinery v0.33.1 ) diff --git a/functions/onprem/orborus/go.sum b/functions/onprem/orborus/go.sum index 3818c38f..a1ec5127 100644 --- a/functions/onprem/orborus/go.sum +++ b/functions/onprem/orborus/go.sum @@ -328,8 +328,8 @@ github.com/sendgrid/sendgrid-go v3.16.1+incompatible h1:zWhTmB0Y8XCDzeWIm2/BIt1G github.com/sendgrid/sendgrid-go v3.16.1+incompatible/go.mod h1:QRQt+LX/NmgVEvmdRw0VT/QgUn499+iza2FnDca9fg8= github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3 h1:n661drycOFuPLCN3Uc8sB6B/s6Z4t2xvBgU1htSHuq8= github.com/sergi/go-diff v1.3.2-0.20230802210424-5b0b94c5c0d3/go.mod h1:A0bzQcvG0E7Rwjx0REVgAGH58e96+X0MeOfepqsbeW4= -github.com/shuffle/shuffle-shared v0.9.15 h1:Gc7c0pbWG6nHWSTkcfAnKgQgWCWfc6aDQ/BIu20z6bM= -github.com/shuffle/shuffle-shared v0.9.15/go.mod h1:PhDEizuz4SmJaSmy0+yrFWwD1mXVUsy8/knKlrqF1qw= +github.com/shuffle/shuffle-shared v0.9.27 h1:YwyWXsp4fCOAPmc1DD+NNf9sVa4RHzp26SvWKxH4ytc= +github.com/shuffle/shuffle-shared v0.9.27/go.mod h1:PhDEizuz4SmJaSmy0+yrFWwD1mXVUsy8/knKlrqF1qw= github.com/sirupsen/logrus v1.7.0/go.mod h1:yWOB1SBYBC5VeMP7gHvWumXLIWorT60ONWic61uBYv0= github.com/sirupsen/logrus v1.9.3 h1:dueUQJ1C2q9oE3F7wvmSGAaVtTmUizReu6fjN8uqzbQ= github.com/sirupsen/logrus v1.9.3/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ= diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index fb93ad89..457b36bc 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -12,7 +12,6 @@ import ( "fmt" "io" "io/ioutil" - "regexp" "log" "math" "net" @@ -20,6 +19,7 @@ import ( "os" "os/exec" "path/filepath" + "regexp" "runtime" "strconv" "strings" @@ -52,6 +52,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" ) @@ -89,6 +90,7 @@ var debug = os.Getenv("DEBUG") == "true" // var baseimagename = "shuffle/shuffle" var baseimageregistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY") var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME") + //var baseimagetagsuffix = os.Getenv("SHUFFLE_BASE_IMAGE_TAG_SUFFIX") // Used for cloud with auth. Onprem in certain cases too. @@ -116,7 +118,7 @@ var pipelineApikey = os.Getenv("SHUFFLE_PIPELINE_AUTH") var pipelineUrl = os.Getenv("SHUFFLE_PIPELINE_URL") var executionIds = []string{} -var pipelines = []shuffle.PipelineInfoMini{} +var pipelines = []shuffle.PipelineInfo{} var namespacemade = false // For K8s var skipPipelineMount = false var tenzirDisabled = false @@ -343,7 +345,7 @@ func deployServiceWorkers(image string) { if len(dockerSwarmBridgeMTU) == 0 { mtu, err = strconv.Atoi(dockerSwarmBridgeMTU) // by default if err != nil { - if debug { + if debug { log.Printf("[DEBUG] Failed to convert the default MTU to int: %s. Using 1500 instead. Input: %s", err, dockerSwarmBridgeMTU) } @@ -530,7 +532,6 @@ func deployServiceWorkers(image string) { nodeCount = uint64(cnt) } - appReplicas := os.Getenv("SHUFFLE_APP_REPLICAS") appReplicaCnt := 2 if len(appReplicas) > 0 { @@ -799,6 +800,48 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { return envVars } +func buildResourcesFromEnv() corev1.ResourceRequirements { + requests := corev1.ResourceList{} + limits := corev1.ResourceList{} + + type item struct { + env string + resourceName corev1.ResourceName + resourceList corev1.ResourceList + } + + items := []item{ + // kubernetes requests + {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", 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 value := strings.TrimSpace(os.Getenv(it.env)); value != "" { + if quantity, err := resource.ParseQuantity(value); err == nil { + it.resourceList[it.resourceName] = quantity + } else { + log.Printf("[WARNING] Cannot parse %s=%q as resource quantity: %v", it.env, value, err) + } + } + } + + rr := corev1.ResourceRequirements{} + if len(requests) > 0 { + rr.Requests = requests + } + if len(limits) > 0 { + rr.Limits = limits + } + + return rr +} + func handleBackendImageDownload(ctx context.Context, images string) error { // Replicate images with lowercase, as the name may be wrong @@ -1042,6 +1085,20 @@ 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"))) + // 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"))) } @@ -1206,6 +1263,7 @@ func deployK8sWorker(image string, identifier string, env []string) error { Image: kubernetesImage, Env: buildEnvVars(envMap), SecurityContext: containerSecurityContext, + Resources: buildResourcesFromEnv(), //ImagePullPolicy: "Never", ImagePullPolicy: corev1.PullIfNotPresent, @@ -1571,12 +1629,12 @@ func initializeImages() { if appSdkVersion == "" { appSdkVersion = "latest" - log.Printf("[WARNING] SHUFFLE_APP_SDK_VERSION not defined. Defaulting to %#v", appSdkVersion) + log.Printf("[INFO] SHUFFLE_APP_SDK_VERSION not defined. Defaulting to %#v", appSdkVersion) } if workerVersion == "" { workerVersion = "latest" - log.Printf("[WARNING] SHUFFLE_WORKER_VERSION not defined. Defaulting to %#v", workerVersion) + log.Printf("[INFO] SHUFFLE_WORKER_VERSION not defined. Defaulting to %#v", workerVersion) } if baseimageregistry == "" { @@ -1586,12 +1644,12 @@ func initializeImages() { if len(os.Getenv("REGISTRY_URL")) > 0 { baseimageregistry = os.Getenv("REGISTRY_URL") } else { - os.Setenv("REGISTRY_URL", baseimageregistry) + // os.Setenv("REGISTRY_URL", baseimageregistry) } os.Setenv("SHUFFLE_BASE_IMAGE_REGISTRY", baseimageregistry) - log.Printf("[WARNING] Setting baseimageregistry to %#v", baseimageregistry) + log.Printf("[INFO] Setting baseimageregistry to %#v", baseimageregistry) } if baseimagename == "" { @@ -1600,7 +1658,7 @@ func initializeImages() { baseimagename = "frikky/shuffle" // Dockerhub os.Setenv("SHUFFLE_BASE_IMAGE_NAME", baseimagename) - log.Printf("[WARNING] Setting baseimagename to %#v", baseimagename) + log.Printf("[INFO] Setting baseimagename to %#v", baseimagename) } // Old sane default overrides: @@ -2063,6 +2121,12 @@ func main() { } } + // Auto enables pipelines IF they are not mentioned + if len(os.Getenv("SHUFFLE_SKIP_PIPELINES")) == 0 { + os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") + os.Setenv("SHUFFLE_PIPELINE_ENABLED", "true") + } + log.Println("[INFO] Setting up execution environment") // //FIXME @@ -2431,9 +2495,10 @@ func main() { for _, incRequest := range executionRequests.Data { // Looking for specific jobs - if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_START" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" { + if incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_START" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" || incRequest.Type == "PIPELINE_UPDATE" { + log.Printf("[INFO] Handling pipeline request from backend: '%s' with argument '%s'", incRequest.Type, incRequest.ExecutionArgument) - os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") + //os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") tenzirDisabled = false // Running NEW or editing pipelines @@ -2456,8 +2521,8 @@ func main() { } else if incRequest.Type == "CATEGORY_UPDATE" { os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") - tenzirDisabled = false + tenzirDisabled = false err = handleFileCategoryChange() if err != nil { log.Printf("[ERROR] Failed to download the file category: %s", err) @@ -2500,7 +2565,7 @@ func main() { log.Printf("[INFO] Got job to start tenzir") // Manual command = overrides to allow starting of Tenzir from the frontend anyway. - os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") + //os.Setenv("SHUFFLE_SKIP_PIPELINES", "false") tenzirDisabled = false // Removed either way @@ -2510,8 +2575,8 @@ func main() { if err != nil { if strings.Contains(fmt.Sprintf("%s", err), "node available") { // Disabling until UI is updated - os.Setenv("SHUFFLE_SKIP_PIPELINES", "true") - tenzirDisabled = true + //os.Setenv("SHUFFLE_SKIP_PIPELINES", "true") + //tenzirDisabled = true log.Printf("[ERROR] Failed to start tenzir, reason: %s", err) err = shuffle.CreateOrgNotification( @@ -2577,7 +2642,7 @@ func main() { executionRequests.Data = executionRequests.Data[0:allowed] } } else if swarmControlMode && (swarmConfig == "run" || swarmConfig == "swarm") { - // any reason it is not maxConcurrency instead of + // any reason it is not maxConcurrency instead of // hardcoded 50? if len(executionRequests.Data) > 50 { executionRequests.Data = executionRequests.Data[0:50] @@ -2759,7 +2824,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { // no need of execution arguments for STOP and DELETE if (incRequest.Type != "PIPELINE_STOP" && incRequest.Type != "PIPELINE_DELETE") && len(incRequest.ExecutionArgument) == 0 { - log.Printf("[ERROR] No execution argument found for pipeline create. Skipping") + log.Printf("[ERROR] No execution argument found for pipeline type %s. Skipping", incRequest.Type) return errors.New("no execution argument found for pipeline create. Skipping") } @@ -2770,6 +2835,7 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { } command := incRequest.ExecutionArgument + pipelines = []shuffle.PipelineInfo{} if incRequest.Type == "PIPELINE_CREATE" { log.Printf("[INFO] Should delete -> recreate new pipeline with id %#v", identifier) //err := deployPipeline(image, identifier, command) @@ -2779,19 +2845,16 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { return err } } else if incRequest.Type == "PIPELINE_DELETE" || incRequest.Type == "PIPELINE_STOP" { - { - log.Printf("[INFO] Should delete pipeline %#v", identifier) - pipelineId, err := searchPipeline(identifier) - if err != nil { - log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) - return err - } + pipelineId := incRequest.ExecutionId + log.Printf("[INFO] Should delete pipeline %#v. PipelineID: %s", identifier, pipelineId) + //pipelineId, err := searchPipeline(identifier) + //if err != nil { + //} - err = deletePipeline(pipelineId) - if err != nil { - log.Printf("[ERROR] Failed Deleting Pipeline %s", err) - return err - } + err = deletePipeline(pipelineId) + if err != nil { + log.Printf("[ERROR] Failed Deleting Pipeline %s", err) + return err } /* @@ -2818,19 +2881,24 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { if err != nil { if err.Error() == "no existing pipeline found with name" { log.Printf("[INFO] Starting a new pipeline with command '%s' and identifier '%s'", command, identifier) - _, CreateErr := createPipeline(command, identifier) - return CreateErr + var createErr error + pipelineId, createErr = createPipeline(command, identifier) + if createErr != nil { + return createErr + } + } else { + log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) + return err } - - log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) - return err } + + log.Printf("[INFO] Starting existing pipeline with ID %s", pipelineId) _, err = updatePipelineState(command, pipelineId, "start") if err != nil { log.Printf("[ERROR] Failed to start Pipeline: %s reason:%s ", pipelineId, err) return err } else { - log.Printf("[INFO] Successfully started the Pipeline: %s", pipelineId) + log.Printf("[INFO] Successfully started pipeline: %s", pipelineId) } } else { @@ -2860,15 +2928,20 @@ func deployTenzirNode() error { ctx := context.Background() cacheKey := "tenzir-key" - - imageName := "frikky/shuffle:tenzir" - containerName := "tenzir-node" - containerStartOptions := container.StartOptions{} _, err = shuffle.GetCache(ctx, cacheKey) if err == nil { return nil } + imageName := "frikky/shuffle:tenzir" + if os.Getenv("TENZIR_IMAGE_NAME") != "" { + imageName = os.Getenv("TENZIR_IMAGE_NAME") + log.Printf("[INFO] Using custom Tenzir image name: %s", imageName) + } + + containerName := "tenzir-node" + containerStartOptions := container.StartOptions{} + containerInfo, err := dockercli.ContainerInspect(ctx, containerName) if err != nil { if dockerclient.IsErrNotFound(err) { @@ -2909,7 +2982,7 @@ func deployTenzirNode() error { } } else { if !containerInfo.State.Running { - log.Printf("[DEBUG] Tenzir Node exists but is not running. Restarting it.") + log.Printf("[DEBUG] Tenzir Node exists, but is not running. Restarting it.") err := dockercli.ContainerStart(ctx, containerName, containerStartOptions) if err != nil { log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) @@ -2959,7 +3032,9 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri ExposedPorts: nat.PortSet{ "5160/tcp": struct{}{}, "514/udp": struct{}{}, + "1514/udp": struct{}{}, "514/tcp": struct{}{}, + "1514/tcp": struct{}{}, }, Entrypoint: []string{containerName}, Env: []string{}, @@ -3007,14 +3082,17 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri PortBindings: nat.PortMap{ "514/tcp": []nat.PortBinding{{HostPort: "514"}}, "514/udp": []nat.PortBinding{{HostPort: "514"}}, + "1514/tcp": []nat.PortBinding{{HostPort: "1514"}}, + "1514/udp": []nat.PortBinding{{HostPort: "1514"}}, "5160/tcp": []nat.PortBinding{{HostPort: "5160"}}, }, Mounts: []mount.Mount{ { Type: "bind", Source: tenzirStorageFolder, - Target: "/var/lib/tenzir/", + Target: "/tmp", }, + /* { Type: "bind", Source: tenzirStorageFolder, @@ -3025,6 +3103,7 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri Source: tenzirStorageFolder, Target: "/var/cache/tenzir/", }, + */ }, VolumeDriver: "local", RestartPolicy: container.RestartPolicy{ @@ -3032,6 +3111,12 @@ func createAndStartTenzirNode(ctx context.Context, containerName, imageName stri }, } + if os.Getenv("SHUFFLE_DISABLE_SYSLOG") == "true" { + hostConfig.PortBindings = nat.PortMap{ + "5160/tcp": []nat.PortBinding{{HostPort: "5160"}}, + } + } + if skipPipelineMount { hostConfig.Mounts = []mount.Mount{} } @@ -3144,7 +3229,7 @@ func createNetworkIfNotExists(ctx context.Context, networkName, subnet, gateway } func checkTenzirNode() error { - if os.Getenv("SHUFFLE_SKIP_PIPELINES") == "true" { + if os.Getenv("SHUFFLE_SKIP_PIPELINES") == "true" && os.Getenv("SHUFFLE_PIPELINE_ENABLED") == "false" { return errors.New("Pipelines are disabled by user with SHUFFLE_SKIP_PIPELINES") } @@ -3217,21 +3302,18 @@ func createPipeline(command, identifier string) (string, error) { //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | sigma /var/lib/tenzir/rule.yaml" //command = "from file /var/lib/tenzir/sysmon_logs.ndjson read json | import" + // Make sure to escape them + //if strings.Contains(command, "/") { + // command = strings.ReplaceAll("\\\"", "", command) + // command = strings.ReplaceAll(command, "\"", "") + //} + requestBody := map[string]interface{}{ - "definition": command, - "name": identifier, - "hidden": false, - "autostart": map[string]bool{ - "created": true, - "completed": false, - "failed": false, - }, - "autodelete": map[string]bool{ - "completed": false, - "failed": false, - "stopped": false, - }, + "definition": command, + "name": identifier, + "hidden": false, "retry_delay": "500.0ms", + "unstoppable": true, } requestBodyJSON, err := json.Marshal(requestBody) @@ -3241,18 +3323,18 @@ func createPipeline(command, identifier string) (string, error) { } forwardData := bytes.NewBuffer(requestBodyJSON) - req, err := http.NewRequest( forwardMethod, url, forwardData, ) + if err != nil { log.Printf("[ERROR] Failed to create HTTP request: %s", err) return "", err } - req.Header.Set("Content-Type", "application/json") + req.Header.Set("Content-Type", "application/json") client := &http.Client{Timeout: 10 * time.Second} resp, err := client.Do(req) if err != nil { @@ -3267,7 +3349,9 @@ func createPipeline(command, identifier string) (string, error) { } if strings.Contains(string(body), "error") { - log.Printf("[ERROR] Pipeline creation response (%d): %s", resp.StatusCode, string(body)) + log.Printf("[ERROR] Pipeline creation error resp (%d): %s", resp.StatusCode, string(body)) + } else { + log.Printf("[DEBUG] Pipeline creation debug (%d): %s", resp.StatusCode, string(body)) } defer resp.Body.Close() @@ -3293,37 +3377,39 @@ func createPipeline(command, identifier string) (string, error) { return "", errors.New("Pipeline ID not found or empty in the response. See error logs.") } - id := response.ID - return id, nil + return response.ID, nil } func updatePipelineState(command, pipelineId, action string) (string, error) { url := fmt.Sprintf("%s/api/v0/pipeline/update", pipelineUrl) forwardMethod := "POST" - requestBody := map[string]interface{}{ "id": pipelineId, - "definition": command, "action": action, + + /* "autostart": map[string]bool{ "created": true, - "completed": true, - "failed": true, + "completed": false, + "failed": false, }, "autodelete": map[string]bool{ "completed": false, "failed": false, "stopped": false, }, + */ } requestBodyJSON, err := json.Marshal(requestBody) if err != nil { return "", err } - forwardData := bytes.NewBuffer(requestBodyJSON) + log.Printf("[INFO] Updating pipeline %s with action %s to ensure it starts. Body: %s", pipelineId, action, string(requestBodyJSON)) + + forwardData := bytes.NewBuffer(requestBodyJSON) req, err := http.NewRequest( forwardMethod, url, @@ -3406,7 +3492,7 @@ func deletePipeline(pipelineId string) error { log.Printf("[INFO] Pipeline with ID: %s deleted successfully", pipelineId) - pipelines = []shuffle.PipelineInfoMini{} + pipelines = []shuffle.PipelineInfo{} return nil } @@ -3480,7 +3566,14 @@ func handleFileCategoryChange() error { } if len(pipelineApikey) == 0 { - return errors.New("Shuffle API-key not set for Pipelines: SHUFFLE_PIPELINE_AUTH=") + //var auth = os.Getenv("AUTH") + //var org = os.Getenv("ORG") + + if len(auth) > 0 && len(org) > 0 { + pipelineApikey = auth + } else { + return errors.New("Shuffle API-key not set for Pipelines: SHUFFLE_PIPELINE_AUTH=") + } } req.Header.Add("Authorization", "Bearer "+pipelineApikey) @@ -3630,6 +3723,7 @@ func removeFileCategory() error { return nil } +// curl https://get.tenzir.app | sh func removeFile(fileName string) error { containerName := "tenzir-node" srcPath := fmt.Sprintf("/var/lib/tenzir/sigma_rules/%s", fileName) @@ -3659,7 +3753,7 @@ func removePath(containerName, path string) error { func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { pipelinePayload := shuffle.LakeConfig{ Enabled: false, - Pipelines: []shuffle.PipelineInfoMini{}, + Pipelines: []shuffle.PipelineInfo{}, } if tenzirDisabled { @@ -3671,18 +3765,9 @@ func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { if len(pipelines) == 0 || randint == 0 { pipelineDef, err := listPipelines() - if err == nil { - for _, pipeline := range pipelineDef { - pipelinePayload.Pipelines = append(pipelinePayload.Pipelines, shuffle.PipelineInfoMini{ - ID: pipeline.ID, - Name: pipeline.Name, - Definition: pipeline.Definition, - TotalRuns: pipeline.TotalRuns, - CreatedAt: pipeline.CreatedAt, - }) - } - - pipelines = pipelinePayload.Pipelines + if err == nil || len(pipelines) > 0 { + pipelines = pipelineDef + pipelinePayload.Pipelines = pipelines } } else { pipelinePayload.Pipelines = pipelines @@ -3695,7 +3780,7 @@ func sendPipelineHealthStatus() (shuffle.LakeConfig, error) { log.Printf("[ERROR] Tenzir node connection problem: %s", err) } else { - tenzirDisabled = true + //tenzirDisabled = true log.Printf("[WARNING] Disabling pipelines: %s. You will need to restart the Orborus to fix this.", err) } @@ -3996,7 +4081,7 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest, image string, // Specific to debugging if len(workerServerUrl) == 0 { - if debug { + if debug { log.Printf("[INFO] Using default worker server url as previous is invalid: %s. Swapping to shuffle-workers:33333", streamUrl) } } diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 81153607..6797ff3c 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -43,6 +43,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" @@ -652,6 +653,7 @@ func deployk8sApp(image string, identifier string, env []string) error { }, }, SecurityContext: containerSecurityContext, + Resources: buildResourcesFromEnv(), }, }, DNSPolicy: corev1.DNSClusterFirst, @@ -2436,6 +2438,49 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { } return envVars } + +func buildResourcesFromEnv() corev1.ResourceRequirements { + requests := corev1.ResourceList{} + limits := corev1.ResourceList{} + + type item struct { + env string + resourceName corev1.ResourceName + resourceList corev1.ResourceList + } + + items := []item{ + // kubernetes requests + {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", 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 value := strings.TrimSpace(os.Getenv(it.env)); value != "" { + if quantity, err := resource.ParseQuantity(value); err == nil { + it.resourceList[it.resourceName] = quantity + } else { + log.Printf("[WARNING] Cannot parse %s=%q as resource quantity: %v", it.env, value, err) + } + } + } + + rr := corev1.ResourceRequirements{} + if len(requests) > 0 { + rr.Requests = requests + } + if len(limits) > 0 { + rr.Limits = limits + } + + return rr +} + func getWorkerBackendExecution(auth string, executionId string) (*shuffle.WorkflowExecution, error) { backendUrl := os.Getenv("BASE_URL") if len(backendUrl) == 0 { @@ -2446,9 +2491,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) @@ -2708,7 +2753,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)) } @@ -2726,7 +2771,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) } @@ -3893,7 +3938,7 @@ func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, if attempts < 2 { // Check the service and fix it. if isKubernetes == "true" { - log.Printf("[WARNING] App Redeployment in K8s isn't fully supported yet, but should be done for app %s with image %s.", appName, image) + log.Printf("[WARNING] App Redeployment in K8s isn't fully supported yet, but should be done for app %s with image %s.", appName, image) } else { _, err = findAppInfo(image, appName, true) if err != nil { @@ -3926,7 +3971,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)) } }