Removed duplicate executions in orborus
This commit is contained in:
@@ -53,7 +53,7 @@ var environment = os.Getenv("ENVIRONMENT_NAME")
|
|||||||
var dockerApiVersion = os.Getenv("DOCKER_API_VERSION")
|
var dockerApiVersion = os.Getenv("DOCKER_API_VERSION")
|
||||||
var runningMode = strings.ToLower(os.Getenv("RUNNING_MODE"))
|
var runningMode = strings.ToLower(os.Getenv("RUNNING_MODE"))
|
||||||
var cleanupEnv = strings.ToLower(os.Getenv("CLEANUP"))
|
var cleanupEnv = strings.ToLower(os.Getenv("CLEANUP"))
|
||||||
var workerIds = []string{}
|
var executionIds = []string{}
|
||||||
|
|
||||||
type ExecutionRequestWrapper struct {
|
type ExecutionRequestWrapper struct {
|
||||||
Data []ExecutionRequest `json:"data"`
|
Data []ExecutionRequest `json:"data"`
|
||||||
@@ -172,7 +172,7 @@ func deployWorker(image string, identifier string, env []string) {
|
|||||||
if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") {
|
if strings.Contains(fmt.Sprintf("%s", err), "Conflict. The container name ") {
|
||||||
uuid := uuid.NewV4()
|
uuid := uuid.NewV4()
|
||||||
identifier = fmt.Sprintf("%s-%s", identifier, uuid)
|
identifier = fmt.Sprintf("%s-%s", identifier, uuid)
|
||||||
log.Printf("2 - Identifier: %s", identifier)
|
log.Printf("[INFO] 2 - Identifier: %s", identifier)
|
||||||
cont, err = dockercli.ContainerCreate(
|
cont, err = dockercli.ContainerCreate(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
config,
|
config,
|
||||||
@@ -221,7 +221,6 @@ func deployWorker(image string, identifier string, env []string) {
|
|||||||
//}
|
//}
|
||||||
} else {
|
} else {
|
||||||
log.Printf("[INFO] Container %s was created under environment %s", cont.ID, environment)
|
log.Printf("[INFO] Container %s was created under environment %s", cont.ID, environment)
|
||||||
//workerIds = append(workerIds, cont.ID)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
return
|
return
|
||||||
@@ -515,8 +514,6 @@ func main() {
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
|
|
||||||
//log.Printf("[INFO] Got %d new requests. Executing: %d. Max: %d", len(executionRequests.Data), executionCount, maxConcurrency)
|
|
||||||
|
|
||||||
allowed := maxConcurrency - executionCount
|
allowed := maxConcurrency - executionCount
|
||||||
if len(executionRequests.Data) > allowed {
|
if len(executionRequests.Data) > allowed {
|
||||||
log.Printf("[WARNING] Throttle - Cutting down requests from %d to %d (MAX: %d, CUR: %d)", len(executionRequests.Data), allowed, maxConcurrency, executionCount)
|
log.Printf("[WARNING] Throttle - Cutting down requests from %d to %d (MAX: %d, CUR: %d)", len(executionRequests.Data), allowed, maxConcurrency, executionCount)
|
||||||
@@ -538,6 +535,23 @@ func main() {
|
|||||||
if execution.Status == "ABORT" || execution.Status == "FAILED" {
|
if execution.Status == "ABORT" || execution.Status == "FAILED" {
|
||||||
log.Printf("[INFO] Executionstatus issue: ", execution.Status)
|
log.Printf("[INFO] Executionstatus issue: ", execution.Status)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
found := false
|
||||||
|
for _, executionId := range executionIds {
|
||||||
|
if execution.ExecutionId == executionId {
|
||||||
|
found = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if found {
|
||||||
|
log.Printf("[INFO] Skipping duplicate %s", execution.ExecutionId)
|
||||||
|
continue
|
||||||
|
} else {
|
||||||
|
//log.Printf("[INFO] Adding to be ran %s", execution.ExecutionId)
|
||||||
|
executionIds = append(executionIds, execution.ExecutionId)
|
||||||
|
}
|
||||||
|
|
||||||
// Now, how do I execute this one?
|
// Now, how do I execute this one?
|
||||||
// FIXME - if error, check the status of the running one. If it's bad, send data back.
|
// FIXME - if error, check the status of the running one. If it's bad, send data back.
|
||||||
containerName := fmt.Sprintf("worker-%s", execution.ExecutionId)
|
containerName := fmt.Sprintf("worker-%s", execution.ExecutionId)
|
||||||
@@ -677,6 +691,7 @@ func getRunningWorkers(ctx context.Context, workerTimeout int) int {
|
|||||||
// FIXME - add this to remove exited workers
|
// FIXME - add this to remove exited workers
|
||||||
// Should it check what happened to the execution? idk
|
// Should it check what happened to the execution? idk
|
||||||
func zombiecheck(ctx context.Context, workerTimeout int) error {
|
func zombiecheck(ctx context.Context, workerTimeout int) error {
|
||||||
|
executionIds = []string{}
|
||||||
log.Println("[INFO] Looking for old containers (zombies)")
|
log.Println("[INFO] Looking for old containers (zombies)")
|
||||||
containers, err := dockercli.ContainerList(ctx, types.ContainerListOptions{
|
containers, err := dockercli.ContainerList(ctx, types.ContainerListOptions{
|
||||||
All: true,
|
All: true,
|
||||||
|
|||||||
Reference in New Issue
Block a user