From abef955448609aa7b5944095fc56d1b8b8c46eb0 Mon Sep 17 00:00:00 2001 From: satti-hari-krishna-reddy Date: Tue, 21 May 2024 14:09:02 +0530 Subject: [PATCH] added image pull checking and removed few log messages --- functions/onprem/orborus/orborus.go | 211 ++++++++++++++++------------ 1 file changed, 125 insertions(+), 86 deletions(-) diff --git a/functions/onprem/orborus/orborus.go b/functions/onprem/orborus/orborus.go index 6a8e260e..dc9e554c 100755 --- a/functions/onprem/orborus/orborus.go +++ b/functions/onprem/orborus/orborus.go @@ -1667,7 +1667,7 @@ func main() { newrequests := []shuffle.ExecutionRequest{} for _, incRequest := range executionRequests.Data { // Looking for specific jobs - if incRequest.Type == "PIPELINE_CREATE" || 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" { err := handlePipeline(incRequest) if err != nil { @@ -2027,40 +2027,38 @@ func main() { // Read from Cache and send it to a webhook // docker run tenzir/tenzir:latest 'from http://192.168.86.44:5002/api/v1/orgs/7e9b9007-5df2-4b47-bca5-c4d267ef2943/cache/CIDR%20ranges?type=text&authorization=cec9d01f-09b2-4419-8a0a-76c6046e3fef read lines | to http://192.168.86.44:5002/api/v1/hooks/webhook_665ace5f-f27b-496a-a365-6e07eb61078c write lines' func handlePipeline(incRequest shuffle.ExecutionRequest) error { + + if tenzirUrl == "" { + tenzirUrl = "http://localhost:5160" + log.Printf("[WARNING] SHUFFLE_TENZIR_URL not set, falling back to default URL: %s",tenzirUrl) + } + err := deployTenzirNode() if err != nil{ log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err) + return err } - // no need of execution arguments for state updates - if incRequest.Type != "PIPELINE_STOP" && len(incRequest.ExecutionArgument) == 0 { + // 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") return errors.New("no execution argument found for pipeline create. Skipping") } - //image := "tenzir/tenzir:latest" identifier := fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-"))) command := incRequest.ExecutionArgument if incRequest.Type == "PIPELINE_CREATE" { - log.Printf("[INFO] Should delete -> recreate new pipeline %#v. Name: %#v", incRequest.ExecutionArgument, identifier) + log.Printf("[INFO] Should delete -> recreate new pipeline with id %#v", identifier) //err := deployPipeline(image, identifier, command) - pipelineId, err := createPipeline(command, identifier) + _, err := createPipeline(command, identifier) if err != nil { log.Printf("[ERROR] Failed to create pipeline: %s", err) return err - } else { - log.Printf("[INFO] Pipeline created successfully with Id: %s", pipelineId) - newErr := savePipelineData(pipelineId, identifier, "running") - if newErr != nil { - log.Printf("[DEBUG] failed to save the pipeline data: %s", newErr) - } else { - log.Printf("[INFO] succesfully saved the pipeline info ") - } } } else if incRequest.Type == "PIPELINE_DELETE" { - log.Printf("[INFO] Should delete pipeline %#v", incRequest.ExecutionArgument) + 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) @@ -2070,6 +2068,8 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { if err != nil { log.Printf("[ERROR] Failed Deleting Pipeline %s", err) return err + } else { + log.Printf("[INFO] successfully deleted the Pipeline: %s", pipelineId) } } else if incRequest.Type == "PIPELINE_STOP" { log.Printf("[INFO] Should stop the pipeline %#v", identifier) @@ -2078,18 +2078,32 @@ func handlePipeline(incRequest shuffle.ExecutionRequest) error { log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) return err } - state, err := updatePipelineState(pipelineId, "stop") + _, err = updatePipelineState(pipelineId, "stop") if err != nil { log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err) return err } else { log.Printf("[INFO] successfully stopped the Pipeline: %s", pipelineId) } - err = savePipelineData(pipelineId, identifier, state) + + } else if incRequest.Type == "PIPELINE_START" { + log.Printf("[INFO] Should start the pipeline %#v", identifier) + pipelineId, err := searchPipeline(identifier) + if err != nil { + if err.Error() == "no existing pipeline found with name" { + log.Printf("[WARNING] no pipeline found for %s, creating a new one", identifier) + _, CreateErr := createPipeline(command, identifier) + return CreateErr + } + log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) + return err + } + _, err = updatePipelineState(pipelineId, "start") if err != nil { - log.Printf("[DEBUG] failed to save the pipeline data: %s", err) + log.Printf("[ERROR] Failed to start Pipeline: %s reason:%s ", pipelineId, err) + return err } else { - log.Printf("[INFO] succesfully saved the pipeline info ") + log.Printf("[INFO] successfully started the Pipeline: %s", pipelineId) } } else { @@ -2106,7 +2120,7 @@ func deployTenzirNode() error { } ctx := context.Background() - cacheKey := "tenzir-key" + cacheKey := "tenzir-key" imageName := "tenzir/tenzir:latest" containerName := "tenzir-node" @@ -2114,19 +2128,29 @@ func deployTenzirNode() error { _, err := shuffle.GetCache(ctx, cacheKey) if err == nil { - return nil - } + return nil + } containerInfo, err := dockercli.ContainerInspect(ctx, containerName) if err != nil { if dockerclient.IsErrNotFound(err) { - pullOptions := types.ImagePullOptions{} - out, err := dockercli.ImagePull(ctx, imageName, pullOptions) - if err != nil { - log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err) + + // Check if image exists + _, _, err := dockercli.ImageInspectWithRaw(ctx, imageName) + if dockerclient.IsErrNotFound(err) { + log.Printf("[DEBUG] pulling image %s", imageName) + pullOptions := types.ImagePullOptions{} + out, err := dockercli.ImagePull(ctx, imageName, pullOptions) + if err != nil { + log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err) + return err + } + defer out.Close() + + io.Copy(io.Discard, out) + } else if err != nil { return err } - defer out.Close() err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions) if err != nil { @@ -2137,37 +2161,35 @@ func deployTenzirNode() error { } } else { if !containerInfo.State.Running { - log.Printf("[DEBUG] Tenzir Node exists but is not running, starting it") + log.Printf("[DEBUG] Tenzir Node exists but is not running") err := dockercli.ContainerStart(ctx, containerName, containerStartOptions) if err != nil { log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) return err } - log.Printf("[INFO] Tenzir Node container started successfully") log.Printf("[INFO] Waiting for Tenzir to become available ...") err = checkTenzirNode() if err != nil { return err } - log.Printf("[INFO] Successfully deployed Tenzir Node!") } } - tenzirStatus := struct { - ContainerStatus string `json:"container_status"` - }{ - ContainerStatus: "running", - } + tenzirStatus := struct { + ContainerStatus string `json:"container_status"` + }{ + ContainerStatus: "running", + } - cacheData, err := json.Marshal(tenzirStatus) - if err != nil { - log.Printf("[WARNING] Failed marshalling execution: %s", err) - } - err = shuffle.SetCache(ctx, cacheKey, cacheData, 1) - if err != nil { + cacheData, err := json.Marshal(tenzirStatus) + if err != nil { + log.Printf("[WARNING] Failed marshalling execution: %s", err) + } + err = shuffle.SetCache(ctx, cacheKey, cacheData, 1) + if err != nil { log.Printf("[WARNING] Failed updating cache for tenzir: %s", err) - } + } return nil } @@ -2265,18 +2287,35 @@ func createPipeline(command, identifier string) (string, error) { toBeDeleted = true } + if strings.Contains(command, "kafka") { + var scheme string + if strings.Contains(command, "http://") { + scheme = "http://" + } else if strings.Contains(command, "https://") { + scheme = "https://" + } + + startIndex := strings.Index(command, scheme) + if startIndex != -1 { + endIndex := startIndex + len(scheme) + endIndex += strings.Index(command[endIndex:], "/") + + command = command[:startIndex] + baseUrl + command[endIndex:] + } + } + requestBody := map[string]interface{}{ "definition": command, "name": identifier, "hidden": false, "autostart": map[string]bool{ "created": true, - "completed": false, - "failed": false, + "completed": true, + "failed": true, }, "autodelete": map[string]bool{ "completed": false, - "failed": true, + "failed": false, "stopped": false, }, "retry_delay": "500.0ms", @@ -2348,12 +2387,12 @@ func updatePipelineState(pipelineId, action string) (string, error) { "action": action, "autostart": map[string]bool{ "created": true, - "completed": false, - "failed": false, + "completed": true, + "failed": true, }, "autodelete": map[string]bool{ "completed": false, - "failed": true, + "failed": false, "stopped": false, }, } @@ -2490,53 +2529,53 @@ func searchPipeline(identifier string) (string, error) { return "", errors.New("no existing pipeline found with name") } -func savePipelineData(pipelineId, identifier, status string) error { +// func savePipelineData(pipelineId, identifier, status string) error { - url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", baseUrl) - identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-") +// url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", baseUrl) +// identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-") - forwardMethod := "PUT" +// forwardMethod := "PUT" - payload := map[string]interface{}{ - "pipeline_id": pipelineId, - "trigger_id": identifierWithoutPrefix, - "status": status, - } +// payload := map[string]interface{}{ +// "pipeline_id": pipelineId, +// "trigger_id": identifierWithoutPrefix, +// "status": status, +// } - payloadBytes, err := json.Marshal(payload) - if err != nil { - log.Printf("[ERROR] Failed to marshal payload: %s", err) - return err - } +// payloadBytes, err := json.Marshal(payload) +// if err != nil { +// log.Printf("[ERROR] Failed to marshal payload: %s", err) +// return err +// } - forwardData := bytes.NewBuffer(payloadBytes) +// forwardData := bytes.NewBuffer(payloadBytes) - 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, 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") - client := &http.Client{Timeout: 10 * time.Second} - resp, err := client.Do(req) - if err != nil { - log.Printf("[ERROR] Failed to send HTTP request: %s", err) - return err - } - defer resp.Body.Close() +// client := &http.Client{Timeout: 10 * time.Second} +// resp, err := client.Do(req) +// if err != nil { +// log.Printf("[ERROR] Failed to send HTTP request: %s", err) +// return err +// } +// defer resp.Body.Close() - if resp.StatusCode != 200 { - log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode) - return fmt.Errorf("unexpected HTTP status code: %d", resp.StatusCode) - } +// if resp.StatusCode != 200 { +// log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode) +// return fmt.Errorf("unexpected HTTP status code: %d", resp.StatusCode) +// } - return nil -} +// return nil +// } // Is this ok to do with Docker? idk :) func getRunningWorkers(ctx context.Context, workerTimeout int) int {