added image pull checking and removed few log messages

This commit is contained in:
satti-hari-krishna-reddy
2024-05-21 14:09:02 +05:30
parent 1f09cce1dd
commit abef955448
+125 -86
View File
@@ -1667,7 +1667,7 @@ func main() {
newrequests := []shuffle.ExecutionRequest{} newrequests := []shuffle.ExecutionRequest{}
for _, incRequest := range executionRequests.Data { for _, incRequest := range executionRequests.Data {
// Looking for specific jobs // 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) err := handlePipeline(incRequest)
if err != nil { if err != nil {
@@ -2027,40 +2027,38 @@ func main() {
// Read from Cache and send it to a webhook // 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' // 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 { 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() err := deployTenzirNode()
if err != nil{ if err != nil{
log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err) log.Printf("[ERROR] failed to deploy the pipeline, reason: %s", err)
return err
} }
// no need of execution arguments for state updates // no need of execution arguments for STOP and DELETE
if incRequest.Type != "PIPELINE_STOP" && len(incRequest.ExecutionArgument) == 0 { 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 create. Skipping")
return errors.New("no execution argument found for pipeline create. Skipping") return errors.New("no execution argument found for pipeline create. Skipping")
} }
//image := "tenzir/tenzir:latest" //image := "tenzir/tenzir:latest"
identifier := fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-"))) identifier := fmt.Sprintf("shuffle-%s", strings.ToLower(strings.ReplaceAll(incRequest.ExecutionSource, " ", "-")))
command := incRequest.ExecutionArgument command := incRequest.ExecutionArgument
if incRequest.Type == "PIPELINE_CREATE" { 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) //err := deployPipeline(image, identifier, command)
pipelineId, err := createPipeline(command, identifier) _, err := createPipeline(command, identifier)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed to create pipeline: %s", err) log.Printf("[ERROR] Failed to create pipeline: %s", err)
return 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" { } 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) pipelineId, err := searchPipeline(identifier)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err) 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 { if err != nil {
log.Printf("[ERROR] Failed Deleting Pipeline %s", err) log.Printf("[ERROR] Failed Deleting Pipeline %s", err)
return err return err
} else {
log.Printf("[INFO] successfully deleted the Pipeline: %s", pipelineId)
} }
} else if incRequest.Type == "PIPELINE_STOP" { } else if incRequest.Type == "PIPELINE_STOP" {
log.Printf("[INFO] Should stop the pipeline %#v", identifier) 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) log.Printf("[ERROR] Failed searching for Pipeline with name %s reason:%s ", identifier, err)
return err return err
} }
state, err := updatePipelineState(pipelineId, "stop") _, err = updatePipelineState(pipelineId, "stop")
if err != nil { if err != nil {
log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err) log.Printf("[ERROR] Failed to stop Pipeline: %s reason:%s ", pipelineId, err)
return err return err
} else { } else {
log.Printf("[INFO] successfully stopped the Pipeline: %s", pipelineId) 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 { 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 { } else {
log.Printf("[INFO] succesfully saved the pipeline info ") log.Printf("[INFO] successfully started the Pipeline: %s", pipelineId)
} }
} else { } else {
@@ -2106,7 +2120,7 @@ func deployTenzirNode() error {
} }
ctx := context.Background() ctx := context.Background()
cacheKey := "tenzir-key" cacheKey := "tenzir-key"
imageName := "tenzir/tenzir:latest" imageName := "tenzir/tenzir:latest"
containerName := "tenzir-node" containerName := "tenzir-node"
@@ -2114,19 +2128,29 @@ func deployTenzirNode() error {
_, err := shuffle.GetCache(ctx, cacheKey) _, err := shuffle.GetCache(ctx, cacheKey)
if err == nil { if err == nil {
return nil return nil
} }
containerInfo, err := dockercli.ContainerInspect(ctx, containerName) containerInfo, err := dockercli.ContainerInspect(ctx, containerName)
if err != nil { if err != nil {
if dockerclient.IsErrNotFound(err) { if dockerclient.IsErrNotFound(err) {
pullOptions := types.ImagePullOptions{}
out, err := dockercli.ImagePull(ctx, imageName, pullOptions) // Check if image exists
if err != nil { _, _, err := dockercli.ImageInspectWithRaw(ctx, imageName)
log.Printf("[ERROR] Failed to pull the Tenzir image: %s", err) 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 return err
} }
defer out.Close()
err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions) err = createAndStartTenzirNode(ctx, containerName, imageName, containerStartOptions)
if err != nil { if err != nil {
@@ -2137,37 +2161,35 @@ func deployTenzirNode() error {
} }
} else { } else {
if !containerInfo.State.Running { 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) err := dockercli.ContainerStart(ctx, containerName, containerStartOptions)
if err != nil { if err != nil {
log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err) log.Printf("[ERROR] Failed to start Tenzir Node container: %v", err)
return err return err
} }
log.Printf("[INFO] Tenzir Node container started successfully")
log.Printf("[INFO] Waiting for Tenzir to become available ...") log.Printf("[INFO] Waiting for Tenzir to become available ...")
err = checkTenzirNode() err = checkTenzirNode()
if err != nil { if err != nil {
return err return err
} }
log.Printf("[INFO] Successfully deployed Tenzir Node!")
} }
} }
tenzirStatus := struct { tenzirStatus := struct {
ContainerStatus string `json:"container_status"` ContainerStatus string `json:"container_status"`
}{ }{
ContainerStatus: "running", ContainerStatus: "running",
} }
cacheData, err := json.Marshal(tenzirStatus) cacheData, err := json.Marshal(tenzirStatus)
if err != nil { if err != nil {
log.Printf("[WARNING] Failed marshalling execution: %s", err) log.Printf("[WARNING] Failed marshalling execution: %s", err)
} }
err = shuffle.SetCache(ctx, cacheKey, cacheData, 1) err = shuffle.SetCache(ctx, cacheKey, cacheData, 1)
if err != nil { if err != nil {
log.Printf("[WARNING] Failed updating cache for tenzir: %s", err) log.Printf("[WARNING] Failed updating cache for tenzir: %s", err)
} }
return nil return nil
} }
@@ -2265,18 +2287,35 @@ func createPipeline(command, identifier string) (string, error) {
toBeDeleted = true 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{}{ requestBody := map[string]interface{}{
"definition": command, "definition": command,
"name": identifier, "name": identifier,
"hidden": false, "hidden": false,
"autostart": map[string]bool{ "autostart": map[string]bool{
"created": true, "created": true,
"completed": false, "completed": true,
"failed": false, "failed": true,
}, },
"autodelete": map[string]bool{ "autodelete": map[string]bool{
"completed": false, "completed": false,
"failed": true, "failed": false,
"stopped": false, "stopped": false,
}, },
"retry_delay": "500.0ms", "retry_delay": "500.0ms",
@@ -2348,12 +2387,12 @@ func updatePipelineState(pipelineId, action string) (string, error) {
"action": action, "action": action,
"autostart": map[string]bool{ "autostart": map[string]bool{
"created": true, "created": true,
"completed": false, "completed": true,
"failed": false, "failed": true,
}, },
"autodelete": map[string]bool{ "autodelete": map[string]bool{
"completed": false, "completed": false,
"failed": true, "failed": false,
"stopped": false, "stopped": false,
}, },
} }
@@ -2490,53 +2529,53 @@ func searchPipeline(identifier string) (string, error) {
return "", errors.New("no existing pipeline found with name") 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) // url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", baseUrl)
identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-") // identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-")
forwardMethod := "PUT" // forwardMethod := "PUT"
payload := map[string]interface{}{ // payload := map[string]interface{}{
"pipeline_id": pipelineId, // "pipeline_id": pipelineId,
"trigger_id": identifierWithoutPrefix, // "trigger_id": identifierWithoutPrefix,
"status": status, // "status": status,
} // }
payloadBytes, err := json.Marshal(payload) // payloadBytes, err := json.Marshal(payload)
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to marshal payload: %s", err) // log.Printf("[ERROR] Failed to marshal payload: %s", err)
return err // return err
} // }
forwardData := bytes.NewBuffer(payloadBytes) // forwardData := bytes.NewBuffer(payloadBytes)
req, err := http.NewRequest( // req, err := http.NewRequest(
forwardMethod, // forwardMethod,
url, // url,
forwardData, // forwardData,
) // )
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err) // log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err // return err
} // }
req.Header.Set("Content-Type", "application/json") // req.Header.Set("Content-Type", "application/json")
client := &http.Client{Timeout: 10 * time.Second} // client := &http.Client{Timeout: 10 * time.Second}
resp, err := client.Do(req) // resp, err := client.Do(req)
if err != nil { // if err != nil {
log.Printf("[ERROR] Failed to send HTTP request: %s", err) // log.Printf("[ERROR] Failed to send HTTP request: %s", err)
return err // return err
} // }
defer resp.Body.Close() // defer resp.Body.Close()
if resp.StatusCode != 200 { // if resp.StatusCode != 200 {
log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode) // log.Printf("[ERROR] Received non-successful HTTP status code: %d", resp.StatusCode)
return fmt.Errorf("unexpected 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 :) // Is this ok to do with Docker? idk :)
func getRunningWorkers(ctx context.Context, workerTimeout int) int { func getRunningWorkers(ctx context.Context, workerTimeout int) int {