trying to improve deployment

This commit is contained in:
Hari Krishna
2024-04-26 16:02:16 +00:00
parent f825306a85
commit c59f575711
2 changed files with 66 additions and 23 deletions
+5
View File
@@ -100,4 +100,9 @@ SHUFFLE_OPENSEARCH_PROXY=
SHUFFLE_OPENSEARCH_INDEX_PREFIX=
SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY=true
#Tenzir related
IS_TENZIR=false
SHUFFLE_TENZIR_URL=http://localhost:5160
DEBUG_MODE=false
+61 -23
View File
@@ -103,6 +103,8 @@ var swarmConfig = os.Getenv("SHUFFLE_SWARM_CONFIG")
var swarmNetworkName = os.Getenv("SHUFFLE_SWARM_NETWORK_NAME")
var orborusLabel = os.Getenv("SHUFFLE_ORBORUS_LABEL")
var memcached = os.Getenv("SHUFFLE_MEMCACHED")
var isTenzir = os.Getenv("IS_TENZIR")
var tenzirUrl = os.Getenv("SHUFFLE_TENZIR_URL")
var executionIds = []string{}
var namespacemade = false // For K8s
@@ -110,6 +112,7 @@ var namespacemade = false // For K8s
var dockercli *dockerclient.Client
var containerId string
var executionCount = 0
var isTenzirReady = false
func init() {
var err error
@@ -1496,13 +1499,17 @@ func main() {
initializeImages()
go func() {
if err := deployTenzirNode(); err != nil {
// Handle the error here
}
}()
if isTenzir == "true" {
func() {
if err := deployTenzirNode(); err != nil {
log.Printf("[ERROR] Failed to deploy the tenzir node, reason: %v", err)
} else {
log.Printf("[INFO] Tenzir node is deployed successfully and is available for requests!")
isTenzirReady = true
}
}()
}
workerImage := fmt.Sprintf("%s/%s/shuffle-worker:%s", baseimageregistry, baseimagename, workerVersion)
if len(newWorkerImage) > 0 {
workerImage = newWorkerImage
@@ -1674,14 +1681,18 @@ 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" {
err := handlePipeline(incRequest)
if err != nil {
log.Printf("[ERROR] Failed handling pipeline: %s", err)
if isTenzir == "true" && incRequest.Type == "PIPELINE_CREATE" || incRequest.Type == "PIPELINE_STOP" || incRequest.Type == "PIPELINE_DELETE" {
if isTenzirReady {
err := handlePipeline(incRequest)
if err != nil {
log.Printf("[ERROR] Failed handling pipeline: %s", err)
//update it to db ??
}
} else {
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
}
log.Printf("[WARNING] Couldnt Handle pipeline request as tenzir node is not ready")
}
toBeRemoved.Data = append(toBeRemoved.Data, incRequest)
} else if incRequest.Type == "DOCKER_IMAGE_DOWNLOAD" {
log.Printf("[INFO] Should delete -> download new image %#v", incRequest.ExecutionArgument)
@@ -2137,12 +2148,6 @@ func deployTenzirNode() error {
Target: "/var/lib/tenzir/",
},
},
LogConfig: container.LogConfig{
Type: "json-file",
Config: map[string]string{
"max-size": "10m",
},
},
VolumeDriver: "local",
}
@@ -2169,7 +2174,7 @@ func deployTenzirNode() error {
log.Printf("[INFO] Existing Tenzir Node container started successfully")
return nil
}
return fmt.Errorf("failed to create Tenzir Node container: %v", err)
return err
}
log.Printf("[INFO] New Tenzir Node container created successfully")
@@ -2181,15 +2186,46 @@ func deployTenzirNode() error {
}
log.Printf("[INFO] New Tenzir Node container started successfully")
log.Printf("[INFO] Waiting for tenzir to become available ...")
err = checkTenzirNode()
if err != nil{
return err
}
return nil
}
func checkTenzirNode() error {
retries := 20
retryInterval := 3 * time.Second
url := fmt.Sprintf("%s/ping",tenzirUrl)
forwardMethod := "POST"
client := http.Client{}
req, err := http.NewRequest(forwardMethod, url, nil)
if err != nil {
log.Printf("[ERROR] Failed to create HTTP request: %s", err)
return err
}
for i := 0; i < retries; i++ {
resp, err := client.Do(req)
if err == nil && resp.StatusCode == http.StatusOK {
return nil
}
time.Sleep(retryInterval)
}
return fmt.Errorf("tenzir node is not available")
}
func createPipeline(command, identifier string) (string, error) {
toBeDeleted := false
pipelineId, err := searchPipeline(identifier)
url := fmt.Sprintf("%s/api/v0/pipelines/create", tenzirUrl)
url := fmt.Sprintf("%s/api/v0/pipeline/create", tenzirUrl)
forwardMethod := "POST"
if err != nil {
@@ -2217,6 +2253,7 @@ func createPipeline(command, identifier string) (string, error) {
"failed": true,
"stopped": false,
},
"retry_delay": "500.0ms",
}
requestBodyJSON, err := json.Marshal(requestBody)
@@ -2429,13 +2466,14 @@ func searchPipeline(identifier string) (string, error) {
func savePipelineData(pipelineId, identifier, status string) error {
url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", tenzirUrl)
url := fmt.Sprintf("%s/api/v1/triggers/pipeline/save", baseUrl)
identifierWithoutPrefix := strings.TrimPrefix(identifier, "shuffle-")
forwardMethod := "PUT"
payload := map[string]interface{}{
"pipeline_id": pipelineId,
"trigger_id": identifier,
"trigger_id": identifierWithoutPrefix,
"status": status,
}