Fixed new caching system in worker to ensure failed apps retry multiple times
This commit is contained in:
@@ -1,5 +1,5 @@
|
|||||||
NAME=shuffle-worker
|
NAME=shuffle-worker
|
||||||
VERSION=0.9.55
|
VERSION=0.9.59
|
||||||
|
|
||||||
echo "Running docker build with $NAME:$VERSION"
|
echo "Running docker build with $NAME:$VERSION"
|
||||||
#CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o worker.bin .
|
#CGO_ENABLED=0 GOOS=linux go build -a -installsuffix cgo -o worker.bin .
|
||||||
|
|||||||
@@ -358,6 +358,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|||||||
log.Printf("\n\n[DEBUG] Result for %s already found - returning\n\n", newExecId)
|
log.Printf("\n\n[DEBUG] Result for %s already found - returning\n\n", newExecId)
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
cacheData := []byte("1")
|
cacheData := []byte("1")
|
||||||
err = shuffle.SetCache(ctx, newExecId, cacheData)
|
err = shuffle.SetCache(ctx, newExecId, cacheData)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -371,17 +372,17 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|||||||
waitTime := time.Duration(action.ExecutionDelay) * time.Second
|
waitTime := time.Duration(action.ExecutionDelay) * time.Second
|
||||||
|
|
||||||
time.AfterFunc(waitTime, func() {
|
time.AfterFunc(waitTime, func() {
|
||||||
DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution)
|
DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
|
||||||
})
|
})
|
||||||
} else {
|
} else {
|
||||||
log.Printf("[DEBUG] Running app %s in docker NORMALLY as there is no delay set", action.Name)
|
log.Printf("[DEBUG] Running app %s in docker NORMALLY as there is no delay set", action.Name)
|
||||||
return DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution)
|
return DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
|
||||||
}
|
}
|
||||||
|
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *container.Config, hostConfig *container.HostConfig, identifier string, workflowExecution shuffle.WorkflowExecution) error {
|
func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *container.Config, hostConfig *container.HostConfig, identifier string, workflowExecution shuffle.WorkflowExecution, newExecId string) error {
|
||||||
cont, err := cli.ContainerCreate(
|
cont, err := cli.ContainerCreate(
|
||||||
ctx,
|
ctx,
|
||||||
config,
|
config,
|
||||||
@@ -396,6 +397,11 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
|
|||||||
if !strings.Contains(err.Error(), "Conflict. The container name") {
|
if !strings.Contains(err.Error(), "Conflict. The container name") {
|
||||||
log.Printf("[ERROR] Container CREATE error (1): %s", err)
|
log.Printf("[ERROR] Container CREATE error (1): %s", err)
|
||||||
|
|
||||||
|
err = shuffle.DeleteCache(ctx, newExecId)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, err)
|
||||||
|
}
|
||||||
|
|
||||||
return err
|
return err
|
||||||
} else {
|
} else {
|
||||||
parsedUuid := uuid.NewV4()
|
parsedUuid := uuid.NewV4()
|
||||||
@@ -414,6 +420,11 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
|
|||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[ERROR] Container create error (2): %s", err)
|
log.Printf("[ERROR] Container create error (2): %s", err)
|
||||||
|
err = shuffle.DeleteCache(ctx, newExecId)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, err)
|
||||||
|
}
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -445,6 +456,12 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
|
|||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[ERROR] Container create error (3): %s", err)
|
log.Printf("[ERROR] Container create error (3): %s", err)
|
||||||
|
|
||||||
|
err = shuffle.DeleteCache(ctx, newExecId)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, err)
|
||||||
|
}
|
||||||
|
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -454,6 +471,12 @@ func DeployContainer(ctx context.Context, cli *dockerclient.Client, config *cont
|
|||||||
|
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[ERROR] Failed to start container in environment %s: %s", environment, err)
|
log.Printf("[ERROR] Failed to start container in environment %s: %s", environment, err)
|
||||||
|
|
||||||
|
err = shuffle.DeleteCache(ctx, newExecId)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] FAILED Deleting cache for %s: %s", newExecId, err)
|
||||||
|
}
|
||||||
|
|
||||||
//shutdown(workflowExecution, workflowExecution.Workflow.ID, true)
|
//shutdown(workflowExecution, workflowExecution.Workflow.ID, true)
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user