From a661741064fe50c344f24ca3463fd24964d521a4 Mon Sep 17 00:00:00 2001 From: frikky Date: Mon, 6 Feb 2023 02:49:53 +0100 Subject: [PATCH] Rebuild worker automatically (CI) --- .env | 2 + backend/app_sdk/app_base.py | 23 ++++++-- backend/go-app/main.go | 4 +- docker-compose.yml | 2 +- functions/onprem/orborus/build.sh | 2 +- functions/onprem/orborus/orborus.go | 90 +++++++++++++++++------------ functions/onprem/worker/build.sh | 5 +- functions/onprem/worker/go.mod | 2 +- functions/onprem/worker/go.sum | 2 + functions/onprem/worker/worker.go | 3 +- 10 files changed, 86 insertions(+), 49 deletions(-) diff --git a/.env b/.env index f5607f7d..687fcf58 100644 --- a/.env +++ b/.env @@ -65,6 +65,8 @@ SHUFFLE_CHAT_DISABLED=false # Controls support chat SHUFFLE_RERUN_SCHEDULE=300 SHUFFLE_DISABLE_RERUN_AND_ABORT=false SHUFFLE_WORKER_SERVER_URL= # Definition in case Worker & Orborus is talking to the wrong server +SHUFFLE_WORKER_SERVER_URL= # Definition in case Worker & Orborus is talking to the wrong server +SHUFFLE_ORBORUS_PULL_TIME= # Definition in case Orborus is pulling too often/not often enough # DATABASE CONFIGURATIONS DATASTORE_EMULATOR_HOST=shuffle-database:8000 diff --git a/backend/app_sdk/app_base.py b/backend/app_sdk/app_base.py index b457a5cf..530a77d9 100644 --- a/backend/app_sdk/app_base.py +++ b/backend/app_sdk/app_base.py @@ -1096,7 +1096,8 @@ class AppBase: get_path = "/api/v1/files/namespaces/%s?execution_id=%s&ids=true" % (category, self.full_execution["execution_id"]) headers = { - "Authorization": "Bearer %s" % self.authorization + "Authorization": "Bearer %s" % self.authorization, + "User-Agent": "Shuffle 1.1.0", } ret = requests.get("%s%s" % (self.url, get_path), headers=headers) @@ -1121,7 +1122,8 @@ class AppBase: get_path = "/api/v1/files/namespaces/%s?execution_id=%s" % (namespace, self.full_execution["execution_id"]) headers = { - "Authorization": "Bearer %s" % self.authorization + "Authorization": "Bearer %s" % self.authorization, + "User-Agent": "Shuffle 1.1.0", } ret1 = requests.get("%s%s" % (self.url, get_path), headers=headers) @@ -1183,7 +1185,8 @@ class AppBase: get_path = "/api/v1/files/%s?execution_id=%s" % (item, full_execution["execution_id"]) headers = { "Content-Type": "application/json", - "Authorization": "Bearer %s" % self.authorization + "Authorization": "Bearer %s" % self.authorization, + "User-Agent": "Shuffle 1.1.0", } ret1 = requests.get("%s%s" % (self.url, get_path), headers=headers) @@ -1285,7 +1288,8 @@ class AppBase: org_id = full_execution["workflow"]["execution_org"]["id"] headers = { "Content-Type": "application/json", - "Authorization": "Bearer %s" % self.authorization + "Authorization": "Bearer %s" % self.authorization, + "User-Agent": "Shuffle 1.1.0", } if not isinstance(infiles, list): @@ -1331,6 +1335,7 @@ class AppBase: new_headers = { "Authorization": f"Bearer {self.authorization}", + "User-Agent": "Shuffle 1.1.0", } upload_path = "/api/v1/files/%s/upload?execution_id=%s" % (cur_id, full_execution["execution_id"]) @@ -1377,7 +1382,8 @@ class AppBase: headers = { "Content-Type": "application/json", - "Authorization": f"Bearer {self.authorization}" + "Authorization": f"Bearer {self.authorization}", + "User-Agent": "Shuffle 1.1.0", } if len(self.action) == 0: @@ -1492,6 +1498,13 @@ class AppBase: self.full_execution = fullexecution + #try: + # if "backend_url" in self.full_execution: + # self.url = self.full_execution["backend_url"] + # self.base_url = self.full_execution["backend_url"] + #except KeyError: + # pass + try: if replace_params == True: for inner_action in self.full_execution["workflow"]["actions"]: diff --git a/backend/go-app/main.go b/backend/go-app/main.go index 49c7b90b..b22bb5f1 100644 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -4302,7 +4302,9 @@ func runInitEs(ctx context.Context) { url := os.Getenv("SHUFFLE_APP_DOWNLOAD_LOCATION") if len(url) == 0 { - url = "https://github.com/frikky/shuffle-apps" + log.Printf("Skipping download since no URL is set") + //url = "https://github.com/frikky/shuffle-apps" + return } username := os.Getenv("SHUFFLE_DOWNLOAD_AUTH_USERNAME") diff --git a/docker-compose.yml b/docker-compose.yml index 4451238d..73881519 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -63,7 +63,7 @@ services: container_name: shuffle-opensearch environment: - bootstrap.memory_lock=true - - "OPENSEARCH_JAVA_OPTS=-Xms1024m -Xmx1024m" # minimum and maximum Java heap size, recommend setting both to 50% of system RAM + - "OPENSEARCH_JAVA_OPTS=-Xms2048m -Xmx2048m" # minimum and maximum Java heap size, recommend setting both to 50% of system RAM - cluster.initial_master_nodes=shuffle-opensearch - cluster.routing.allocation.disk.threshold_enabled=false - cluster.name=shuffle-cluster diff --git a/functions/onprem/orborus/build.sh b/functions/onprem/orborus/build.sh index 26265cda..42ac3769 100644 --- a/functions/onprem/orborus/build.sh +++ b/functions/onprem/orborus/build.sh @@ -1,5 +1,5 @@ NAME=shuffle-orborus -VERSION=1.1.1 +VERSION=1.1.2 echo "Running docker build with $NAME:$VERSION" #docker rmi frikky/shuffle:$NAME --force diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 3c71db51..39363d19 100644 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -7,7 +7,7 @@ package main // FIXME: // 2022/01/12 17:13:36 [WARNING] Swarm init: Error response from daemon: manager stopped: failed to listen on remote API address: listen tcp: address tcp/2377%!(EXTRA string=172.23.0.2): unknown port -// frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/frikky/shuffle-worker:nightly +// frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/shuffle/shuffle-worker:nightly // Potential issues: // Default network could be same as on the host @@ -60,9 +60,9 @@ var appSdkVersion = os.Getenv("SHUFFLE_APP_SDK_VERSION") var workerVersion = os.Getenv("SHUFFLE_WORKER_VERSION") var newWorkerImage = os.Getenv("SHUFFLE_WORKER_IMAGE") -// var baseimagename = "docker.pkg.github.com/frikky/shuffle" +// var baseimagename = "docker.pkg.github.com/shuffle/shuffle" // var baseimagename = "ghcr.io/frikky" -// var baseimagename = "frikky/shuffle" +// var baseimagename = "shuffle/shuffle" var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME") var baseimageregistry = os.Getenv("SHUFFLE_BASE_IMAGE_REGISTRY") var baseimagetagsuffix = os.Getenv("SHUFFLE_BASE_IMAGE_TAG_SUFFIX") @@ -82,6 +82,8 @@ var timezone = os.Getenv("TZ") var containerName = os.Getenv("ORBORUS_CONTAINER_NAME") var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG") var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME") +var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL") + var executionIds = []string{} var dockercli *dockerclient.Client @@ -207,7 +209,7 @@ func deployServiceWorkers(image string) { ctx := context.Background() // Looks for and cleans up all existing items in swarm we can't re-use (Shuffle only) - // frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/frikky/shuffle-worker:nightly + // frikky@debian:~/git/shuffle/functions/onprem/worker$ docker service create --replicas 5 --name shuffle-workers --env SHUFFLE_SWARM_CONFIG=run --publish published=33333,target=33333 ghcr.io/shuffle/shuffle-worker:nightly networkName := "shuffle_swarm_executions" if len(swarmNetworkName) > 0 { networkName = swarmNetworkName @@ -591,7 +593,7 @@ func deployWorker(image string, identifier string, env []string, executionReques 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) + //log.Printf("[INFO] 2 - Identifier: %s", identifier) cont, err = dockercli.ContainerCreate( context.Background(), config, @@ -714,6 +716,7 @@ func initializeImages() { baseimageregistry = "ghcr.io" log.Printf("[DEBUG] Setting baseimageregistry") } + if baseimagename == "" { baseimagename = "shuffle/shuffle" // Dockerhub baseimagename = "shuffle" // Github @@ -877,6 +880,15 @@ func main() { timezone = "Europe/Amsterdam" } + if len(os.Getenv("SHUFFLE_ORBORUS_PULL_TIME")) > 0 { + log.Printf("[INFO] Trying to set Orborus sleep time between polls to %s", os.Getenv("SHUFFLE_ORBORUS_PULL_TIME")) + + tmpInt, err := strconv.Atoi(os.Getenv("SHUFFLE_ORBORUS_PULL_TIME")) + if err == nil { + sleepTime = tmpInt + } + } + log.Printf("[INFO] Running with timezone %s", timezone) workerTimeout := 600 @@ -996,6 +1008,11 @@ func main() { req.Header.Add("Org", org) } + if len(orborusLabel) > 0 { + log.Printf("[DEBUG] Sending with Label %s", orborusLabel) + req.Header.Add("X-Orborus-Label", orborusLabel) + } + log.Printf("[INFO] Waiting for executions at %s with Environment %#v", fullUrl, environment) hasStarted := false for { @@ -1015,22 +1032,9 @@ func main() { continue } - // FIXME - add check for StatusCode - if newresp.StatusCode != 200 { - if hasStarted { - log.Printf("[WARNING] Bad statuscode from backend: %d", newresp.StatusCode) - } - } else { - if !hasStarted { - log.Printf("[DEBUG] Starting iteration. Got statuscode %d from backend on first request", newresp.StatusCode) - } - - hasStarted = true - } - body, err := ioutil.ReadAll(newresp.Body) if err != nil { - log.Printf("[ERROR] Failed reading body: %s", err) + log.Printf("[ERROR] Failed reading body from Shuffle: %s", err) zombiecounter += 1 if zombiecounter*sleepTime > workerTimeout { go zombiecheck(ctx, workerTimeout) @@ -1040,6 +1044,17 @@ func main() { continue } + // FIXME - add check for StatusCode + if newresp.StatusCode != 200 { + log.Printf("[ERROR] Backend configuration missing (%d): %s", newresp.StatusCode, string(body)) + } else { + if !hasStarted { + log.Printf("[DEBUG] Starting iteration. Got statuscode %d from backend on first request", newresp.StatusCode) + } + + hasStarted = true + } + var executionRequests shuffle.ExecutionRequestWrapper err = json.Unmarshal(body, &executionRequests) if err != nil { @@ -1093,11 +1108,11 @@ func main() { var toBeRemoved shuffle.ExecutionRequestWrapper for _, execution := range executionRequests.Data { if len(execution.ExecutionArgument) > 0 { - log.Printf("[INFO] Argument: %#v", execution.ExecutionArgument) + log.Printf("[INFO] Argument: %s", execution.ExecutionArgument) } if execution.Type == "schedule" { - log.Printf("[INFO] SOMETHING ELSE :O: %s", execution.Type) + log.Printf("[INFO] Schedule type! Weird deployment. Type: %s", execution.Type) continue } @@ -1105,23 +1120,10 @@ func main() { log.Printf("[INFO] Executionstatus issue: ", execution.Status) } - /* - found := false - for _, executionId := range executionIds { - if execution.ExecutionId == executionId { - found = true - break - } - } - - // Doesn't work because of USER INPUT - if found { - log.Printf("[INFO] Skipping duplicate %s", execution.ExecutionId) - continue - } else { - //log.Printf("[INFO] Adding to be ran %s", execution.ExecutionId) - } - */ + if shuffle.ArrayContains(executionIds, execution.ExecutionId) { + log.Printf("[INFO] Execution already handled: %s", execution.ExecutionId) + continue + } // Now, how do I execute this one? // FIXME - if error, check the status of the running one. If it's bad, send data back. @@ -1190,6 +1192,18 @@ func main() { result.Header.Add("Content-Type", "application/json") result.Header.Add("Org-Id", environment) + if len(auth) > 0 { + result.Header.Add("Authorization", auth) + } + + if len(org) > 0 { + result.Header.Add("Org", org) + } + + if len(orborusLabel) > 0 { + result.Header.Add("X-Orborus-Label", orborusLabel) + } + resultResp, err := client.Do(result) if err != nil { log.Printf("[ERROR] Failed making confirm request: %s", err) diff --git a/functions/onprem/worker/build.sh b/functions/onprem/worker/build.sh index 25fc3778..1d69dab9 100644 --- a/functions/onprem/worker/build.sh +++ b/functions/onprem/worker/build.sh @@ -1,5 +1,5 @@ NAME=shuffle-worker -VERSION=1.1.0 +VERSION=1.1.1 echo "Running docker build with $NAME:$VERSION" #CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o worker.bin . @@ -13,3 +13,6 @@ docker build . -t frikky/shuffle:$NAME -t frikky/shuffle:$NAME_$VERSION -t docke #docker tag frikky/shuffle:$NAME ghcr.io/frikky/shuffle-worker:0.8.52 docker push ghcr.io/frikky/$NAME:$VERSION docker push ghcr.io/frikky/$NAME:nightly + +docker push ghcr.io/shuffle/$NAME:$VERSION +docker push ghcr.io/shuffle/$NAME:nightly diff --git a/functions/onprem/worker/go.mod b/functions/onprem/worker/go.mod index 0808095a..40bbd225 100644 --- a/functions/onprem/worker/go.mod +++ b/functions/onprem/worker/go.mod @@ -10,6 +10,6 @@ require ( github.com/docker/go-connections v0.4.0 // indirect github.com/gorilla/mux v1.8.0 github.com/patrickmn/go-cache v2.1.0+incompatible - github.com/shuffle/shuffle-shared v0.3.24 + github.com/shuffle/shuffle-shared v0.3.51 go4.org v0.0.0-20201209231011-d4a079459e60 // indirect ) diff --git a/functions/onprem/worker/go.sum b/functions/onprem/worker/go.sum index 02ba7801..1e874329 100644 --- a/functions/onprem/worker/go.sum +++ b/functions/onprem/worker/go.sum @@ -818,6 +818,8 @@ github.com/shuffle/shuffle-shared v0.2.27 h1:YT9MtXyMSxIGMpNovjp9pCKFyt2gk40EdAX github.com/shuffle/shuffle-shared v0.2.27/go.mod h1:YuMle0RjwXb3hxR5PdaOOD9e+hUyK34OABS0UbrT/Sk= github.com/shuffle/shuffle-shared v0.3.24 h1:zBDZan4u2XjC6TAi5BdFoVroBPGYd6PAha+3/cSfD6w= github.com/shuffle/shuffle-shared v0.3.24/go.mod h1:yI6HCog/R3Kq1FvCIVbXedLl87rtSuDOyzolmuMswB4= +github.com/shuffle/shuffle-shared v0.3.51 h1:+JPEGw6R4a320who+SrGP/VqBxZdKPdcLw/07cO7d6c= +github.com/shuffle/shuffle-shared v0.3.51/go.mod h1:jQrYySmvp/0De5ftrAaY6xwwr7TMfqBmBxQ2AX9yrjQ= github.com/shurcooL/sanitized_anchor_name v1.0.0/go.mod h1:1NzhyTcUVG4SuEtjjoZeVRXNmyL/1OwPU0+IJeTBvfc= github.com/sirupsen/logrus v1.0.4-0.20170822132746-89742aefa4b2/go.mod h1:pMByvHTf9Beacp5x1UXfOR9xyW/9antXMhjMPG0dEzc= github.com/sirupsen/logrus v1.0.6/go.mod h1:pMByvHTf9Beacp5x1UXfOR9xyW/9antXMhjMPG0dEzc= diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index c32c26cf..6810b971 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -252,7 +252,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] } cacheData := []byte("1") - err = shuffle.SetCache(ctx, newExecId, cacheData) + err = shuffle.SetCache(ctx, newExecId, cacheData, 30) if err != nil { log.Printf("[WARNING] Failed setting cache for action %s: %s", newExecId, err) } else { @@ -1614,6 +1614,7 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec } log.Printf("[INFO] Status: %s, Results: %d, actions: %d", workflowExecution.Status, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions)+extra) + if workflowExecution.Status != "EXECUTING" { log.Printf("[WARNING] Exiting as worker execution has status %s!", workflowExecution.Status) log.Printf("[DEBUG] Shutting down (21)")