From 5226f795519bec0eb3f3be83d8de8374176d8f33 Mon Sep 17 00:00:00 2001 From: Frikky Date: Mon, 30 Sep 2024 13:41:43 +0200 Subject: [PATCH] Made Orborus and Worker both able to download images properly. Added a timeout to make sure images get the time to be built. This is a 30 second stupid addition --- functions/onprem/orborus/go.mod | 2 +- functions/onprem/orborus/orborus.go | 63 +++++++++------ functions/onprem/worker/worker.go | 118 ++-------------------------- 3 files changed, 48 insertions(+), 135 deletions(-) diff --git a/functions/onprem/orborus/go.mod b/functions/onprem/orborus/go.mod index 433e9fd2..33c8a332 100644 --- a/functions/onprem/orborus/go.mod +++ b/functions/onprem/orborus/go.mod @@ -4,7 +4,7 @@ go 1.22.0 toolchain go1.22.2 -//replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared +replace github.com/shuffle/shuffle-shared => ../../../../shuffle-shared require ( github.com/docker/docker v27.0.2+incompatible diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index d6be28d6..7572e806 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -114,6 +114,8 @@ var dockercli *dockerclient.Client var containerId string var executionCount = 0 +var imagedownloadTimeout = time.Second * 300 + func init() { var err error @@ -722,34 +724,46 @@ func buildEnvVars(envMap map[string]string) []corev1.EnvVar { } func handleBackendImageDownload(ctx context.Context, images string) error { - // Should use docker to: - // 1. Pull the image & tag it - // 2. Distribute the image by updating service if "run" + + // Replicate images with lowercase, as the name may be wrong + // Most of the time lowercase is correct. Swapping to have that first + originalImages := images + images = strings.ToLower(images) + "," + originalImages + + log.Printf("[DEBUG] Should remove existing image (s): %s", images) + + // Remove the image + removeOptions := image.RemoveOptions{} + for _, image := range strings.Split(images, ",") { + image = strings.TrimSpace(image) + if !strings.Contains(image, "/") { + image = fmt.Sprintf("frikky/shuffle:%s", image) + } + + // There is no real point in actual removal. This may however be a good idea, as Worker will force download the new one anyway + resp, err := dockercli.ImageRemove(ctx, image, removeOptions) + if err != nil { + log.Printf("[ERROR] Failed removing image: %s", err) + } else { + log.Printf("[DEBUG] Removed image: %s", resp) + } + + + err = shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) + if err != nil { + log.Printf("[ERROR] Failed downloading image: %s", err) + } else { + log.Printf("[DEBUG] Downloaded image: %s", image) + //break + } + } + if swarmConfig == "run" || swarmConfig == "swarm" { log.Printf("[DEBUG] Should update service with new image after updating(s): %s. \n\nNOT IMPLEMENTED: Contact support@shuffler.io for support.\n\n", images) // 1. Download the image // 2. Find the existing service using the image // 3. Update the service with the new image in a rolling restart - } else { - log.Printf("[DEBUG] Should remove existing image (s): %s", images) - - // Remove the image - removeOptions := image.RemoveOptions{} - - for _, image := range strings.Split(images, ",") { - image = strings.TrimSpace(image) - if !strings.Contains(image, "/") { - image = fmt.Sprintf("frikky/shuffle:%s", image) - } - - resp, err := dockercli.ImageRemove(ctx, image, removeOptions) - if err != nil { - log.Printf("[ERROR] Failed removing image: %s", err) - } else { - log.Printf("[DEBUG] Removed image: %s", resp) - } - } } return nil @@ -2036,10 +2050,15 @@ func main() { log.Printf("[INFO] Should delete -> download new image %#v", incRequest.ExecutionArgument) if len(incRequest.ExecutionArgument) > 0 { + // FIXME: Wait X seconds before running this as the image build may not be done yet. This is shitty, but may be ok to do in Orborus. Easy fix for the future: Just let it run through jobs 5-10 times before actually picking it up + time.Sleep(time.Duration(25) * time.Second) + err = handleBackendImageDownload(ctx, incRequest.ExecutionArgument) if err != nil { log.Printf("[ERROR] Failed handling image delete -> download: %s", err) } + } else { + log.Printf("[ERROR] No image name provided for download. Removing job from queue.") } toBeRemoved.Data = append(toBeRemoved.Data, incRequest) diff --git a/functions/onprem/worker/worker.go b/functions/onprem/worker/worker.go index 2625ce80..3402ea6f 100644 --- a/functions/onprem/worker/worker.go +++ b/functions/onprem/worker/worker.go @@ -769,10 +769,10 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env [] // image on every Orborus/new worker restart. // Running as coroutine for eventual completeness - //go downloadDockerImageBackend(&http.Client{}, image) + //go shuffle.DownloadDockerImageBackend(&http.Client{}, image) // FIXME: With goroutines it got too much trouble of deploying with an older version // Allowing slow startups, as long as it's eventually fast, and uses the same registry as on host. - downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) + shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) } var exposedPort int @@ -1499,7 +1499,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { return } - err := downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) + err := shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) executed := false if err == nil { log.Printf("[DEBUG] Downloaded image %s from backend (CLEANUP)", image) @@ -1612,7 +1612,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) { } log.Printf("[DEBUG][%s] Failed deploy. Downloading image %s: %s", workflowExecution.ExecutionId, image, err) - err := downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) + err := shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image) executed := false if err == nil { @@ -2902,113 +2902,7 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener { return listener } -func downloadDockerImageBackend(client *http.Client, imageName string) error { - // Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD - if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" { - log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. Not downloading image %s", imageName) - return nil - } - if arrayContains(downloadedImages, imageName) { - log.Printf("[DEBUG] Image %s already downloaded - not re-downloading", imageName) - return nil - } - - log.Printf("[DEBUG] Trying to download image %s from backend %s as it doesn't exist. All images: %#v", imageName, baseUrl, downloadedImages) - - downloadedImages = append(downloadedImages, imageName) - - data := fmt.Sprintf(`{"name": "%s"}`, imageName) - dockerImgUrl := fmt.Sprintf("%s/api/v1/get_docker_image", baseUrl) - - req, err := http.NewRequest( - "POST", - dockerImgUrl, - bytes.NewBuffer([]byte(data)), - ) - - authorization := os.Getenv("AUTHORIZATION") - if len(authorization) > 0 { - req.Header.Add("Authorization", fmt.Sprintf("Bearer %s", authorization)) - } else { - log.Printf("[WARNING] No auth found - running backend download without it.") - //return - } - - newresp, err := topClient.Do(req) - if err != nil { - log.Printf("[ERROR] Failed download request for %s: %s", imageName, err) - return err - } - - defer newresp.Body.Close() - if newresp.StatusCode != 200 { - log.Printf("[ERROR] Docker download for image %s (backend) StatusCode (1): %d", imageName, newresp.StatusCode) - return errors.New(fmt.Sprintf("Failed to get image - status code %d", newresp.StatusCode)) - } - - newImageName := strings.Replace(imageName, "/", "_", -1) - newFileName := newImageName + ".tar" - - tar, err := os.Create(newFileName) - if err != nil { - log.Printf("[WARNING] Failed creating file: %s", err) - return err - } - - defer tar.Close() - _, err = io.Copy(tar, newresp.Body) - if err != nil { - log.Printf("[WARNING] Failed response body copying: %s", err) - return err - } - tar.Seek(0, 0) - - dockercli, err := dockerclient.NewEnvClient() - if err != nil { - log.Printf("[ERROR] Unable to create docker client (3): %s", err) - return err - } - - defer dockercli.Close() - - imageLoadResponse, err := dockercli.ImageLoad(context.Background(), tar, true) - if err != nil { - log.Printf("[ERROR] Error loading images: %s", err) - return err - } - - defer imageLoadResponse.Body.Close() - body, err := ioutil.ReadAll(imageLoadResponse.Body) - if err != nil { - log.Printf("[ERROR] Error reading: %s", err) - return err - } - - if strings.Contains(string(body), "no such file") { - return errors.New(string(body)) - } - - baseTag := strings.Split(imageName, ":") - if len(baseTag) > 1 { - tag := baseTag[1] - log.Printf("[DEBUG] Creating tag copies of downloaded containers from tag %s", tag) - - // Remapping - ctx := context.Background() - dockercli.ImageTag(ctx, imageName, fmt.Sprintf("frikky/shuffle:%s", tag)) - dockercli.ImageTag(ctx, imageName, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", tag)) - - downloadedImages = append(downloadedImages, fmt.Sprintf("frikky/shuffle:%s", tag)) - downloadedImages = append(downloadedImages, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", tag)) - - } - - os.Remove(newFileName) - - log.Printf("[INFO] Successfully loaded image %s: %s", imageName, string(body)) - return nil -} func findActiveSwarmNodes(dockercli *dockerclient.Client) (int64, error) { ctx := context.Background() @@ -4159,7 +4053,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) { } log.Printf("[INFO] Downloading image %s", image.Image) - downloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image.Image) + shuffle.DownloadDockerImageBackend(&http.Client{Timeout: imagedownloadTimeout}, image.Image) // return success resp.WriteHeader(200) @@ -4213,4 +4107,4 @@ func runWebserver(listener net.Listener) { if err != nil { log.Printf("[ERROR] Serve issue in worker: %#v", err) } -} \ No newline at end of file +}