diff --git a/backend/app_sdk/Dockerfile b/backend/app_sdk/Dockerfile index ac322bb9..d3709032 100644 --- a/backend/app_sdk/Dockerfile +++ b/backend/app_sdk/Dockerfile @@ -1,6 +1,6 @@ #FROM python:3.9.1-alpine as base -#FROM python:3.10.0-alpine as base -FROM python:3.11.3-alpine as base +FROM python:3.10.0-alpine as base +#FROM python:3.11.3-alpine as base FROM base as builder RUN apk --no-cache add --update alpine-sdk libffi libffi-dev musl-dev openssl-dev tzdata coreutils diff --git a/backend/app_sdk/app_base.py b/backend/app_sdk/app_base.py index b4ab4ff7..39073de7 100644 --- a/backend/app_sdk/app_base.py +++ b/backend/app_sdk/app_base.py @@ -2279,13 +2279,6 @@ class AppBase: # Can't handle self yet (?) ret = run.render(**globals()) - - # Load output as JSON - try: - ret = json.loads(ret) - except: - pass - return ret except jinja2.exceptions.TemplateNotFound as e: self.logger.info(f"[ERROR] Liquid Template error: {e}") @@ -3071,6 +3064,7 @@ class AppBase: #self.logger.info(action["parameters"]) # This seems redundant now + self.logger.info("[DEBUG] Pre parameters") for parameter in newparams: action["parameters"].append(parameter) @@ -3092,6 +3086,7 @@ class AppBase: # Multi_parameter has the data for each. variable minlength = 0 + self.logger.info("[DEBUG] Pre-loading parameters") multi_parameters = json.loads(json.dumps(params)) multiexecution = False multi_execution_lists = [] @@ -3518,11 +3513,8 @@ class AppBase: try: del params[field] self.logger.info("[WARNING] Removed field invalid field %s" % field) - except KeyError as e: - self.logger.info("[WARNING] Tried to remove field %s but it didn't exist" % field) + except KeyError: break - else: - self.logger.info("[ERROR] Couldn't find fieldsplit in error. Raw error: %s" % errorstring) else: newres = json.dumps({ "success": False, @@ -3899,11 +3891,7 @@ class AppBase: else: self.logger.info("ACTION TYPE (unhandled): %s" % type(action)) - #await app.execute_action(app.action) app.execute_action(app.action) - #app.run(host="0.0.0.0", port=33334) - if __name__ == "__main__": AppBase.run() - #asyncio.run(AppBase.run(), debug=True) diff --git a/backend/app_sdk/build.sh b/backend/app_sdk/build.sh index 1ff45c75..dd540b94 100644 --- a/backend/app_sdk/build.sh +++ b/backend/app_sdk/build.sh @@ -2,7 +2,7 @@ ### DEFAULT NAME=shuffle-app_sdk -VERSION=1.1.0 +VERSION=1.2.0 docker rmi docker.pkg.github.com/frikky/shuffle/$NAME:$VERSION --force docker build . -f Dockerfile -t frikky/shuffle:app_sdk -t frikky/$NAME:$VERSION -t docker.pkg.github.com/frikky/shuffle/$NAME:$VERSION -t ghcr.io/frikky/$NAME:$VERSION -t ghcr.io/frikky/$NAME:nightly -t shuffle/shuffle:app_sdk -t shuffle/$NAME:$VERSION -t docker.pkg.github.com/shuffle/shuffle/$NAME:$VERSION -t ghcr.io/shuffle/$NAME:$VERSION -t ghcr.io/shuffle/$NAME:nightly diff --git a/backend/go-app/docker.go b/backend/go-app/docker.go index dbf1958e..00202dbd 100644 --- a/backend/go-app/docker.go +++ b/backend/go-app/docker.go @@ -678,8 +678,11 @@ func getDockerImage(resp http.ResponseWriter, request *http.Request) { alternativeName = strings.Join(alternativeNameSplit[1:3], "/") } + log.Printf("[INFO] Trying to download image: %s. Alt: %s", version.Name, alternativeName) + for _, image := range images { for _, tag := range image.RepoTags { + //log.Printf("[DEBUG] Tag: %s", tag) if strings.ToLower(tag) == strings.ToLower(version.Name) { img = image tagFound = tag @@ -693,6 +696,29 @@ func getDockerImage(resp http.ResponseWriter, request *http.Request) { } } + pullOptions := types.ImagePullOptions{} + if len(img.ID) == 0 { + _, err := dockercli.ImagePull(context.Background(), version.Name, pullOptions) + if err == nil { + tagFound = version.Name + img.ID = version.Name + img2.ID = version.Name + + dockercli.ImageTag(ctx, version.Name, alternativeName) + } + } + + if len(img2.ID) == 0 { + _, err := dockercli.ImagePull(context.Background(), alternativeName, pullOptions) + if err == nil { + tagFound = alternativeName + img.ID = alternativeName + img2.ID = alternativeName + + dockercli.ImageTag(ctx, alternativeName, version.Name) + } + } + // REBUILDS THE APP if len(img.ID) == 0 { if len(img2.ID) == 0 { @@ -722,7 +748,7 @@ func getDockerImage(resp http.ResponseWriter, request *http.Request) { foundApp := shuffle.WorkflowApp{} imageName = strings.ToLower(imageName) imageVersion = strings.ToLower(imageVersion) - log.Printf("[DEBUG] Looking for appname %s with version %s", imageName, imageVersion) + log.Printf("[DEBUG] Docker Looking for appname %s with version %s", imageName, imageVersion) for _, app := range workflowapps { if strings.ToLower(strings.Replace(app.Name, " ", "_", -1)) == imageName && app.AppVersion == imageVersion { diff --git a/backend/go-app/go.mod b/backend/go-app/go.mod index 95ba41de..69662589 100644 --- a/backend/go-app/go.mod +++ b/backend/go-app/go.mod @@ -19,7 +19,7 @@ require ( github.com/gorilla/mux v1.8.0 github.com/h2non/filetype v1.1.3 github.com/satori/go.uuid v1.2.0 - github.com/shuffle/shuffle-shared v0.4.2 + github.com/shuffle/shuffle-shared v0.4.9 golang.org/x/crypto v0.3.0 google.golang.org/api v0.103.0 google.golang.org/appengine v1.6.7 diff --git a/backend/go-app/main.go b/backend/go-app/main.go index 08756aad..dc8059b3 100644 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -81,18 +81,14 @@ import ( var gceProject = "shuffle" var bucketName = "shuffler.appspot.com" var baseAppPath = "/home/frikky/git/shaffuru/tmp/apps" + var baseDockerName = "frikky/shuffle" var registryName = "registry.hub.docker.com" var runningEnvironment = "onprem" var syncUrl = "https://shuffler.io" - -// var syncUrl = "http://localhost:5002" var syncSubUrl = "https://shuffler.io" -//var syncUrl = "http://localhost:5002" -//var syncSubUrl = "https://050196912a9d.ngrok.io" - var dbclient *datastore.Client type Userapi struct { @@ -5873,7 +5869,7 @@ func initHandlers() { r.HandleFunc("/api/v1/streams/results", handleGetStreamResults).Methods("POST", "OPTIONS") // Used by orborus - r.HandleFunc("/api/v1/workflows/queue", handleGetWorkflowqueue).Methods("GET") + r.HandleFunc("/api/v1/workflows/queue", handleGetWorkflowqueue).Methods("GET", "POST") r.HandleFunc("/api/v1/workflows/queue/confirm", handleGetWorkflowqueueConfirm).Methods("POST") // App specific @@ -6018,6 +6014,8 @@ func initHandlers() { r.HandleFunc("/api/v1/users/notifications/clear", shuffle.HandleClearNotifications).Methods("GET", "OPTIONS") r.HandleFunc("/api/v1/users/notifications/{notificationId}/markasread", shuffle.HandleMarkAsRead).Methods("GET", "OPTIONS") + r.HandleFunc("/api/v1/conversation", shuffle.RunActionAI).Methods("POST", "OPTIONS") + //r.HandleFunc("/api/v1/users/notifications/{notificationId}/markasread", shuffle.HandleMarkAsRead).Methods("GET", "OPTIONS") r.HandleFunc("/api/v1/dashboards/{key}/widgets", shuffle.HandleNewWidget).Methods("POST", "OPTIONS") r.HandleFunc("/api/v1/dashboards/{key}/widgets/{widget_id}", shuffle.HandleGetWidget).Methods("GET", "OPTIONS") diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index d410a3ed..91bc6dbd 100644 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -11,6 +11,7 @@ import ( "io" "io/ioutil" "log" + "math/rand" "net/http" "net/url" "os" @@ -251,16 +252,42 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { return } - id := request.Header.Get("Org-Id") - if len(id) == 0 { - log.Printf("[INFO] No org-id header set") + // This is really the environment's name - NOT org-id + orgId := request.Header.Get("Org-Id") + if len(orgId) == 0 { + log.Printf("[AUDIT] No org-id header set") resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Specify the org-id header."}`))) return } - ctx := context.Background() - env, err := shuffle.GetEnvironment(ctx, id, "") + environment := request.Header.Get("org") + if len(environment) == 0 { + log.Printf("[AUDIT] No 'org' header set (get workflow queue). Required for cloud.") + /* + resp.WriteHeader(403) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Specify the org header. This can be done by setting the 'ORG' environment variable for Orborus to your Org ID in Shuffle"}`))) + return + */ + } + + orborusLabel := request.Header.Get("x-orborus-label") + + // This section is cloud custom for now + auth := request.Header.Get("Authorization") + if len(auth) == 0 { + log.Printf("[AUDIT] No Authorization header set. Required for cloud. Env: %s, org: %s", orgId, environment) + /* + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Specify the auth header (only applicable for cloud for now)."}`))) + return + */ + } + + //log.Printf("[AUDIT] Get workflow queue for org %s, env %s, orborus label %s", orgId, environment, orborusLabel) + + ctx := shuffle.GetContext(request) + env, err := shuffle.GetEnvironment(ctx, orgId, "") timeNow := time.Now().Unix() if err == nil && len(env.Id) > 0 && len(env.Name) > 0 { if time.Now().Unix() > env.Edited+60 { @@ -273,7 +300,154 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { } } - executionRequests, err := shuffle.GetWorkflowQueue(ctx, id, 100) + //log.Printf("Found env: %#v", env) + if len(env.OrgId) > 0 { + environment = env.OrgId + } + + if request.Method == "POST" { + if rand.Intn(1) == 0 { + // Parse out body + body, err := ioutil.ReadAll(request.Body) + if err == nil { + + // Parse out CPU, memory and disk. + + var envData shuffle.OrborusStats + err = json.Unmarshal(body, &envData) + if err == nil && !envData.Swarm && !envData.Kubernetes && (envData.CPU > 0 || envData.Memory > 0 || envData.Disk > 0) { + + // Set the input in memory + envData.OrgId = orgId + envData.Environment = environment + envData.OrborusLabel = orborusLabel + envData.Timestamp = time.Now().Unix() + + if envData.CPU > 0 && envData.MaxCPU > 0 { + envData.CPUPercent = float64(envData.CPU) / float64(envData.MaxCPU) + } + + if envData.Memory > 0 && envData.MaxMemory > 0 { + envData.MemoryPercent = float64(envData.Memory) / float64(envData.MaxMemory) + } + + // Check if CPU percent constantly has stayed above X% for the last Y requests + percentageCheck := 90 + concurrentChecks := 0 + + //if int(envData.CPUPercent) > percentageCheck { + // Get cached data + percentages := []float64{} + cacheKey := fmt.Sprintf("%s_%s_percent", orgId, strings.ToLower(environment)) + + // Marshal float list into []byte + cacheData := []byte{} + cache, err := shuffle.GetCache(ctx, cacheKey) + if err == nil { + // Unmarshal into percentages + cacheData := []byte(cache.([]uint8)) + err = json.Unmarshal(cacheData, &percentages) + if err != nil { + log.Printf("[INFO] error in cache unmarshal for percentages: %s", err) + } + + if len(percentages) > concurrentChecks { + percentages = percentages[:concurrentChecks] + } + + percentages = append(percentages, envData.CPUPercent) + if len(percentages) > concurrentChecks { + //log.Printf("[INFO] Checking percentages: %v", percentages) + + // percentageCheck := 1 + sendAlert := true + for _, p := range percentages { + if int(p) < percentageCheck { + //log.Printf("[AUDIT] CPU percent is below %d: %d", percentageCheck, int(p)) + sendAlert = false + break + } + } + + if sendAlert { + log.Printf("[INFO] CPU percent has been above %d percent for the last 5 requests. Sending alert. Env: %s, org: %s", percentageCheck, environment, orgId) + + // Set notification + alert for organization + err = shuffle.CreateOrgNotification( + ctx, + fmt.Sprintf("CPU percent has been above %d percent", percentageCheck), + fmt.Sprintf("A environment %s has been using more than %d\\% CPU for the last 5 requests.", environment, percentageCheck), + fmt.Sprintf("/admin?tab=environments"), + environment, + true, + ) + + if err != nil { + log.Printf("[ERROR] error creating notification: %s", err) + } + + org, err := shuffle.GetOrg(ctx, environment) + if err == nil { + foundRecommendation := false + for _, recommendation := range org.Priorities { + if strings.Contains(recommendation.Name, "CPU") { + foundRecommendation = true + break + } + } + + if !foundRecommendation { + // Add to start of org.Priorities + org.Priorities = append(org.Priorities, shuffle.Priority{ + Name: fmt.Sprintf("High CPU in environment %s", orgId), + Description: fmt.Sprintf("The environment %s has been using more than %d percent CPU.", orgId, percentageCheck), + Type: "scale", + Active: true, + URL: fmt.Sprintf("/admin?tab=environments"), + }) + + //Make last item the first item + org.Priorities = append([]shuffle.Priority{org.Priorities[len(org.Priorities)-1]}, org.Priorities[:len(org.Priorities)-1]...) + err = shuffle.SetOrg(ctx, *org, org.Id) + if err != nil { + log.Printf("[ERROR] Problem setting org: %s", err) + } + } + } + } + + if len(percentages) > 1 { + percentages = percentages[1:] + } + } + + // Marshal float list into []byte + } else { + //log.Printf("[ERROR] Failed getting cache: %s", err) + percentages = append(percentages, envData.CPUPercent) + } + + if len(percentages) > 0 { + //log.Printf("[DEBUG] Setting cache for %s: %#v", cacheKey, percentages) + cacheData, err = json.Marshal(percentages) + if err != nil { + log.Printf("[INFO] error in cache marshal: %s", err) + } + + // Add the new data + go shuffle.SetCache(ctx, cacheKey, cacheData, 5) + } + } + + //log.Printf("CPU percent: %f", envData.CPUPercent) + //log.Printf("Memory percent: %f", envData.MemoryPercent*100) + + go shuffle.SetenvStats(ctx, envData) + } + } + } + + executionRequests, err := shuffle.GetWorkflowQueue(ctx, orgId, 100) if err != nil { // Skipping as this comes up over and over //log.Printf("(2) Failed reading body for workflowqueue: %s", err) @@ -290,21 +464,21 @@ func handleGetWorkflowqueue(resp http.ResponseWriter, request *http.Request) { // Try again :) if len(env.Id) == 0 && len(env.Name) == 0 { - orgId := "" + foundId := "" for _, requestData := range executionRequests.Data { execution, err := shuffle.GetWorkflowExecution(ctx, requestData.ExecutionId) if err == nil { if len(execution.ExecutionOrg) > 0 { - orgId = execution.ExecutionOrg + foundId = execution.ExecutionOrg break } } } if len(orgId) > 0 { - env, err := shuffle.GetEnvironment(ctx, id, orgId) + env, err := shuffle.GetEnvironment(ctx, orgId, foundId) if err != nil { - log.Printf("[WARNING] No env found matching %s - continuing without updating orborus anyway: %s", id, err) + log.Printf("[WARNING] No env found matching %s - continuing without updating orborus anyway: %s", orgId, err) //resp.WriteHeader(401) //resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "No env found matching %s"}`, id))) //return @@ -361,9 +535,9 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) { err = json.Unmarshal(body, &actionResult) if err != nil { log.Printf("[WARNING] Failed ActionResult unmarshaling (stream result): %s", err) - resp.WriteHeader(401) - resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) - return + //resp.WriteHeader(401) + //resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + //return } ctx := context.Background() @@ -476,9 +650,9 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) { err = json.Unmarshal(body, &actionResult) if err != nil { log.Printf("[WARNING] Failed ActionResult unmarshaling (queue): %s", err) - resp.WriteHeader(401) - resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) - return + //resp.WriteHeader(401) + //resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) + //return } //log.Printf("Received action: %#v", actionResult) diff --git a/docker-compose.yml b/docker-compose.yml index 1143f01c..b6878df7 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -1,62 +1,62 @@ version: '3' services: - frontend: - image: ghcr.io/shuffle/shuffle-frontend:latest - container_name: shuffle-frontend - hostname: shuffle-frontend - ports: - - "${FRONTEND_PORT}:80" - - "${FRONTEND_PORT_HTTPS}:443" - networks: - - shuffle - environment: - - BACKEND_HOSTNAME=${BACKEND_HOSTNAME} - restart: unless-stopped - depends_on: - - backend - backend: - image: ghcr.io/shuffle/shuffle-backend:latest - container_name: shuffle-backend - hostname: ${BACKEND_HOSTNAME} - # Here for debugging: - ports: - - "${BACKEND_PORT}:5001" - networks: - - shuffle - volumes: - - /var/run/docker.sock:/var/run/docker.sock - - ${SHUFFLE_APP_HOTLOAD_LOCATION}:/shuffle-apps:z - - ${SHUFFLE_FILE_LOCATION}:/shuffle-files:z - env_file: .env - environment: - #- DOCKER_HOST=tcp://docker-socket-proxy:2375 - - SHUFFLE_APP_HOTLOAD_FOLDER=/shuffle-apps - - SHUFFLE_FILE_LOCATION=/shuffle-files - restart: unless-stopped - orborus: - image: ghcr.io/shuffle/shuffle-orborus:latest - container_name: shuffle-orborus - hostname: shuffle-orborus - networks: - - shuffle - volumes: - - /var/run/docker.sock:/var/run/docker.sock - environment: - #- DOCKER_HOST=tcp://docker-socket-proxy:2375 - - SHUFFLE_WORKER_VERSION=latest - - ENVIRONMENT_NAME=${ENVIRONMENT_NAME} - - BASE_URL=http://${OUTER_HOSTNAME}:5001 - - DOCKER_API_VERSION=1.40 - - SHUFFLE_BASE_IMAGE_NAME=${SHUFFLE_BASE_IMAGE_NAME} - - SHUFFLE_BASE_IMAGE_REGISTRY=${SHUFFLE_BASE_IMAGE_REGISTRY} - - SHUFFLE_BASE_IMAGE_TAG_SUFFIX=${SHUFFLE_BASE_IMAGE_TAG_SUFFIX} - - HTTP_PROXY=${HTTP_PROXY} - - HTTPS_PROXY=${HTTPS_PROXY} - - SHUFFLE_PASS_WORKER_PROXY=${SHUFFLE_PASS_WORKER_PROXY} - - SHUFFLE_PASS_APP_PROXY=${SHUFFLE_PASS_APP_PROXY} - restart: unless-stopped - security_opt: - - seccomp:unconfined + #frontend: + # image: ghcr.io/shuffle/shuffle-frontend:latest + # container_name: shuffle-frontend + # hostname: shuffle-frontend + # ports: + # - "${FRONTEND_PORT}:80" + # - "${FRONTEND_PORT_HTTPS}:443" + # networks: + # - shuffle + # environment: + # - BACKEND_HOSTNAME=${BACKEND_HOSTNAME} + # restart: unless-stopped + # depends_on: + # - backend + #backend: + # image: ghcr.io/shuffle/shuffle-backend:latest + # container_name: shuffle-backend + # hostname: ${BACKEND_HOSTNAME} + # # Here for debugging: + # ports: + # - "${BACKEND_PORT}:5001" + # networks: + # - shuffle + # volumes: + # - /var/run/docker.sock:/var/run/docker.sock + # - ${SHUFFLE_APP_HOTLOAD_LOCATION}:/shuffle-apps:z + # - ${SHUFFLE_FILE_LOCATION}:/shuffle-files:z + # env_file: .env + # environment: + # #- DOCKER_HOST=tcp://docker-socket-proxy:2375 + # - SHUFFLE_APP_HOTLOAD_FOLDER=/shuffle-apps + # - SHUFFLE_FILE_LOCATION=/shuffle-files + # restart: unless-stopped + #orborus: + # image: ghcr.io/shuffle/shuffle-orborus:latest + # container_name: shuffle-orborus + # hostname: shuffle-orborus + # networks: + # - shuffle + # volumes: + # - /var/run/docker.sock:/var/run/docker.sock + # environment: + # #- DOCKER_HOST=tcp://docker-socket-proxy:2375 + # - SHUFFLE_WORKER_VERSION=latest + # - ENVIRONMENT_NAME=${ENVIRONMENT_NAME} + # - BASE_URL=http://${OUTER_HOSTNAME}:5001 + # - DOCKER_API_VERSION=1.40 + # - SHUFFLE_BASE_IMAGE_NAME=${SHUFFLE_BASE_IMAGE_NAME} + # - SHUFFLE_BASE_IMAGE_REGISTRY=${SHUFFLE_BASE_IMAGE_REGISTRY} + # - SHUFFLE_BASE_IMAGE_TAG_SUFFIX=${SHUFFLE_BASE_IMAGE_TAG_SUFFIX} + # - HTTP_PROXY=${HTTP_PROXY} + # - HTTPS_PROXY=${HTTPS_PROXY} + # - SHUFFLE_PASS_WORKER_PROXY=${SHUFFLE_PASS_WORKER_PROXY} + # - SHUFFLE_PASS_APP_PROXY=${SHUFFLE_PASS_APP_PROXY} + # restart: unless-stopped + # security_opt: + # - seccomp:unconfined opensearch: image: opensearchproject/opensearch:2.5.0 hostname: shuffle-opensearch diff --git a/frontend/public/images/experienced.png b/frontend/public/images/experienced.png new file mode 100644 index 00000000..8e5f3e5a Binary files /dev/null and b/frontend/public/images/experienced.png differ diff --git a/frontend/public/images/logos/orange_logo.svg b/frontend/public/images/logos/orange_logo.svg new file mode 100644 index 00000000..1024dd6a --- /dev/null +++ b/frontend/public/images/logos/orange_logo.svg @@ -0,0 +1,5 @@ + diff --git a/frontend/public/images/welcome-to-shuffle.png b/frontend/public/images/welcome-to-shuffle.png new file mode 100644 index 00000000..cb976e31 Binary files /dev/null and b/frontend/public/images/welcome-to-shuffle.png differ diff --git a/frontend/src/App.jsx b/frontend/src/App.jsx index 06c91245..41922ff4 100644 --- a/frontend/src/App.jsx +++ b/frontend/src/App.jsx @@ -10,7 +10,7 @@ import GettingStarted from "./views/GettingStarted"; import EditWebhook from "./views/EditWebhook"; import AngularWorkflow from "./views/AngularWorkflow"; -import Header from "./components/Header"; +import Header from "./components/Header.jsx"; import theme from "./theme"; import Apps from "./views/Apps"; import AppCreator from "./views/AppCreator"; @@ -336,7 +336,9 @@ const App = (message, props) => { userdata={userdata} {...props} /> + {/*
+ */}