Merge branch 'main' into 1.4.0
This commit is contained in:
@@ -19,6 +19,7 @@ import (
|
||||
"io"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"math"
|
||||
"net"
|
||||
"net/http"
|
||||
"os"
|
||||
@@ -26,9 +27,8 @@ import (
|
||||
"runtime"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
"sync"
|
||||
"math"
|
||||
"time"
|
||||
|
||||
//"os/signal"
|
||||
//"syscall"
|
||||
@@ -63,8 +63,8 @@ import (
|
||||
var sleepTime = 2
|
||||
|
||||
// Making it work on low-end machines even during busy times :)
|
||||
// May cause some things to run slowly
|
||||
var maxConcurrency = 7
|
||||
// May cause some things to run slowly
|
||||
var maxConcurrency = 7
|
||||
|
||||
// Timeout if something rashes
|
||||
var workerTimeoutEnv = os.Getenv("SHUFFLE_ORBORUS_EXECUTION_TIMEOUT")
|
||||
@@ -76,7 +76,8 @@ var dockerSwarmBridgeMTU = os.Getenv("SHUFFLE_SWARM_BRIDGE_DEFAULT_MTU")
|
||||
var dockerSwarmBridgeInterface = os.Getenv("SHUFFLE_SWARM_BRIDGE_DEFAULT_INTERFACE")
|
||||
var isKubernetes = os.Getenv("IS_KUBERNETES")
|
||||
var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
|
||||
var maxCPUPercent = 95
|
||||
var maxCPUPercent = 90
|
||||
|
||||
|
||||
// var baseimagename = "docker.pkg.github.com/shuffle/shuffle"
|
||||
// var baseimagename = "ghcr.io/frikky"
|
||||
@@ -482,6 +483,7 @@ func deployServiceWorkers(image string) {
|
||||
fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")),
|
||||
fmt.Sprintf("SHUFFLE_MAX_SWARM_NODES=%d", os.Getenv("SHUFFLE_MAX_SWARM_NODES")),
|
||||
fmt.Sprintf("SHUFFLE_BASE_IMAGE_NAME=%s", os.Getenv("SHUFFLE_BASE_IMAGE_NAME")),
|
||||
fmt.Sprintf("SHUFFLE_APP_REQUEST_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_REQUEST_TIMEOUT")),
|
||||
},
|
||||
//Hosts: []string{
|
||||
// innerContainerName,
|
||||
@@ -781,9 +783,9 @@ func deployWorker(image string, identifier string, env []string, executionReques
|
||||
|
||||
log.Printf("[INFO] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace)
|
||||
return nil
|
||||
}
|
||||
}
|
||||
|
||||
// Binds is the actual "-v" volume.
|
||||
// Binds is the actual "-v" volume.
|
||||
// Max 20% CPU every second
|
||||
|
||||
//CPUQuota: 25000,
|
||||
@@ -1138,7 +1140,7 @@ func getOrborusStats(ctx context.Context) shuffle.OrborusStats {
|
||||
newStats.MaxQueue = maxConcurrency
|
||||
newStats.Queue = executionCount
|
||||
|
||||
if isKubernetes == "true" || runningMode == "kubernetes" || runningMode == "k8s" {
|
||||
if isKubernetes == "true" || runningMode == "kubernetes" || runningMode == "k8s" {
|
||||
newStats.Kubernetes = true
|
||||
return newStats
|
||||
}
|
||||
@@ -1171,6 +1173,7 @@ func getOrborusStats(ctx context.Context) shuffle.OrborusStats {
|
||||
|
||||
// Get list of all running containers
|
||||
containers, err := dockercli.ContainerList(ctx, container.ListOptions{})
|
||||
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed getting container list: %s", err)
|
||||
return newStats
|
||||
@@ -1229,33 +1232,33 @@ func getOrborusStats(ctx context.Context) shuffle.OrborusStats {
|
||||
// check if it's NaN or Inf
|
||||
if !math.IsNaN(result.cpuUsage) {
|
||||
totalCPU += float64(result.cpuUsage)
|
||||
}
|
||||
}
|
||||
|
||||
if !math.IsNaN(result.memoryUsage) {
|
||||
memUsage += float64(result.memoryUsage)
|
||||
}
|
||||
}
|
||||
|
||||
newStats.CPUPercent = totalCPU/float64(newStats.CPU)
|
||||
newStats.CPUPercent = totalCPU / float64(newStats.CPU)
|
||||
newStats.MemoryPercent = memUsage
|
||||
|
||||
//log.Printf("[DEBUG] CPU: %.2f, Memory: %.2f", newStats.CPUPercent, newStats.MemoryPercent)
|
||||
|
||||
/*
|
||||
cpuPercent, err := cpu.Percent(250*time.Millisecond, false)
|
||||
if err == nil && len(cpuPercent) > 0 {
|
||||
newStats.CPUPercent = cpuPercent[0]
|
||||
}
|
||||
//Percent(interval time.Duration, percpu bool) ([]float64, error)
|
||||
cpuPercent, err := cpu.Percent(250*time.Millisecond, false)
|
||||
if err == nil && len(cpuPercent) > 0 {
|
||||
newStats.CPUPercent = cpuPercent[0]
|
||||
}
|
||||
//Percent(interval time.Duration, percpu bool) ([]float64, error)
|
||||
|
||||
// Get memory usage
|
||||
memory, err := memory.Get()
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed getting memory stats: %s", err)
|
||||
} else {
|
||||
newStats.Memory = int(memory.Used)
|
||||
newStats.MaxMemory = int(memory.Total)
|
||||
}
|
||||
// Get memory usage
|
||||
memory, err := memory.Get()
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Failed getting memory stats: %s", err)
|
||||
} else {
|
||||
newStats.Memory = int(memory.Used)
|
||||
newStats.MaxMemory = int(memory.Total)
|
||||
}
|
||||
*/
|
||||
|
||||
// Get disk usage
|
||||
@@ -1560,13 +1563,20 @@ func main() {
|
||||
}
|
||||
|
||||
if os.Getenv("SHUFFLE_MAX_CPU") != "" {
|
||||
// parse
|
||||
// parse
|
||||
tmpInt, err := strconv.Atoi(os.Getenv("SHUFFLE_MAX_CPU"))
|
||||
if err == nil {
|
||||
maxCPUPercent = tmpInt
|
||||
}
|
||||
}
|
||||
|
||||
swarmPollingTime := time.Now()
|
||||
swarmRequestsMade := 0
|
||||
swarmControlMode := false
|
||||
if os.Getenv("SHUFFLE_SWARM_CONTROL_MODE") == "true" {
|
||||
swarmControlMode = true
|
||||
}
|
||||
|
||||
log.Printf("[INFO] Waiting for executions at %s with Environment %#v", fullUrl, environment)
|
||||
hasStarted := false
|
||||
for {
|
||||
@@ -1732,6 +1742,20 @@ func main() {
|
||||
log.Printf("[WARNING] Throttle - Cutting down requests from %d to %d (MAX: %d, CUR: %d)", len(executionRequests.Data), allowed, maxConcurrency, executionCount)
|
||||
executionRequests.Data = executionRequests.Data[0:allowed]
|
||||
}
|
||||
} else if (swarmControlMode && (swarmConfig == "run" || swarmConfig == "swarm")) {
|
||||
if len(executionRequests.Data) > 50 {
|
||||
executionRequests.Data = executionRequests.Data[0:50]
|
||||
}
|
||||
|
||||
if swarmRequestsMade > 100 && time.Since(swarmPollingTime).Seconds() > 5 {
|
||||
log.Printf("[DEBUG] Swarm requests made: %d", swarmRequestsMade)
|
||||
time.Sleep(time.Duration(sleepTime) * time.Second)
|
||||
|
||||
swarmPollingTime = time.Now()
|
||||
swarmRequestsMade = 0
|
||||
}
|
||||
|
||||
swarmRequestsMade += len(executionRequests.Data)
|
||||
}
|
||||
|
||||
// New, abortable version. Should check executionid and remove everything else
|
||||
@@ -1760,9 +1784,9 @@ func main() {
|
||||
|
||||
// Should check when last this was ran, and if it's more than 10 minutes ago and it's not finished, we should run it again?
|
||||
/*
|
||||
if swarmConfig != "run" && swarmConfig != "swarm" {
|
||||
continue
|
||||
}
|
||||
if swarmConfig != "run" && swarmConfig != "swarm" {
|
||||
continue
|
||||
}
|
||||
*/
|
||||
}
|
||||
|
||||
@@ -1818,7 +1842,6 @@ func main() {
|
||||
env = append(env, fmt.Sprintf("SHUFFLE_VOLUME_BINDS=%s", os.Getenv("SHUFFLE_VOLUME_BINDS")))
|
||||
}
|
||||
|
||||
|
||||
if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 {
|
||||
env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
|
||||
}
|
||||
@@ -2585,6 +2608,7 @@ func getRunningWorkers(ctx context.Context, workerTimeout int) int {
|
||||
if isKubernetes == "true" {
|
||||
log.Printf("[INFO] Getting running workers in kubernetes")
|
||||
|
||||
|
||||
thresholdTime := time.Now().Add(time.Duration(-workerTimeout) * time.Second)
|
||||
|
||||
clientset, _, err := getKubernetesClient()
|
||||
@@ -2616,20 +2640,20 @@ func getRunningWorkers(ctx context.Context, workerTimeout int) int {
|
||||
containers, err := dockercli.ContainerList(ctx, container.ListOptions{
|
||||
All: true,
|
||||
})
|
||||
|
||||
|
||||
// Automatically updates the version
|
||||
if err != nil {
|
||||
log.Printf("[ERROR] Error getting containers: %s", err)
|
||||
|
||||
|
||||
newVersionSplit := strings.Split(fmt.Sprintf("%s", err), "version is")
|
||||
if len(newVersionSplit) > 1 {
|
||||
//dockerApiVersion = strings.TrimSpace(newVersionSplit[1])
|
||||
log.Printf("[DEBUG] WANT to change the API version to default to %s?", strings.TrimSpace(newVersionSplit[1]))
|
||||
}
|
||||
|
||||
|
||||
return maxConcurrency
|
||||
}
|
||||
|
||||
|
||||
currenttime := time.Now().Unix()
|
||||
|
||||
for _, container := range containers {
|
||||
@@ -2642,7 +2666,7 @@ func getRunningWorkers(ctx context.Context, workerTimeout int) int {
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// Check image name
|
||||
if !shuffleFound {
|
||||
continue
|
||||
@@ -2650,13 +2674,13 @@ func getRunningWorkers(ctx context.Context, workerTimeout int) int {
|
||||
//} else {
|
||||
// log.Printf("NAME: %s", container.Image)
|
||||
}
|
||||
|
||||
|
||||
for _, name := range container.Names {
|
||||
// FIXME - add name_version_uid_uid regex check as well
|
||||
if !strings.HasPrefix(name, "/worker") {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
//log.Printf("Time: %d - %d", currenttime-container.Created, int64(workerTimeout))
|
||||
if container.State == "running" && currenttime-container.Created < int64(workerTimeout) {
|
||||
counter += 1
|
||||
@@ -2800,7 +2824,7 @@ func sendWorkerRequest(workflowExecution shuffle.ExecutionRequest) error {
|
||||
streamUrl = fmt.Sprintf("%s:33333/api/v1/execute", parsedBaseurl)
|
||||
}
|
||||
|
||||
if len(workerServerUrl) > 0 {
|
||||
if len(workerServerUrl) > 0 {
|
||||
streamUrl = fmt.Sprintf("%s:33333/api/v1/execute", workerServerUrl)
|
||||
}
|
||||
|
||||
|
||||
@@ -3,7 +3,6 @@ package main
|
||||
import (
|
||||
"github.com/shuffle/shuffle-shared"
|
||||
|
||||
|
||||
"bytes"
|
||||
"context"
|
||||
"encoding/json"
|
||||
@@ -22,8 +21,8 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/docker/docker/api/types"
|
||||
"github.com/docker/docker/api/types/filters"
|
||||
"github.com/docker/docker/api/types/container"
|
||||
"github.com/docker/docker/api/types/filters"
|
||||
"github.com/docker/docker/api/types/mount"
|
||||
dockerclient "github.com/docker/docker/client"
|
||||
// This is for automatic removal of certain code :)
|
||||
@@ -57,6 +56,7 @@ var kubernetesNamespace = os.Getenv("KUBERNETES_NAMESPACE")
|
||||
|
||||
// var baseimagename = os.Getenv("SHUFFLE_BASE_IMAGE_NAME")
|
||||
|
||||
|
||||
// var baseimagename = "registry.hub.docker.com/frikky/shuffle"
|
||||
var registryName = "registry.hub.docker.com"
|
||||
var sleepTime = 2
|
||||
@@ -81,6 +81,7 @@ var startAction string
|
||||
//var allLogs map[string]string
|
||||
//var containerIds []string
|
||||
var downloadedImages []string
|
||||
|
||||
type ImageDownloadBody struct {
|
||||
Image string `json:"image"`
|
||||
}
|
||||
@@ -92,7 +93,6 @@ type ImageRequest struct {
|
||||
var finishedExecutions []string
|
||||
var imagesDistributed []string
|
||||
|
||||
|
||||
// Images to be autodeployed in the latest version of Shuffle.
|
||||
var autoDeploy = map[string]string{
|
||||
"http:1.4.0": "frikky/shuffle:http_1.4.0",
|
||||
@@ -139,7 +139,6 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
return err
|
||||
}
|
||||
|
||||
|
||||
handleExecutionResult(workflowExecution)
|
||||
validated := shuffle.ValidateFinished(ctx, -1, workflowExecution)
|
||||
if validated {
|
||||
@@ -179,7 +178,7 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
}
|
||||
}
|
||||
|
||||
if len(subflowId) == 0 {
|
||||
if len(subflowId) == 0 {
|
||||
log.Printf("[DEBUG][%s] No waiting result found. Not polling", workflowExecution.ExecutionId)
|
||||
|
||||
for _, action := range workflowExecution.Workflow.Actions {
|
||||
@@ -187,19 +186,17 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
workflowExecution.Workflow.Triggers = append(workflowExecution.Workflow.Triggers, shuffle.Trigger{
|
||||
AppName: action.AppName,
|
||||
Parameters: action.Parameters,
|
||||
ID: action.ID,
|
||||
ID: action.ID,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
for _, trigger := range workflowExecution.Workflow.Triggers {
|
||||
//log.Printf("[DEBUG] Found trigger %s", trigger.AppName)
|
||||
if trigger.AppName != "User Input" && trigger.AppName != "Shuffle Workflow" && trigger.AppName != "shuffle-subflow" {
|
||||
continue
|
||||
}
|
||||
|
||||
|
||||
// check if it has wait for results in params
|
||||
wait := false
|
||||
for _, param := range trigger.Parameters {
|
||||
@@ -218,9 +215,9 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
//log.Printf("[DEBUG][%s] Found result %s", workflowExecution.ExecutionId, result.Action.ID)
|
||||
if result.Action.ID == trigger.ID && result.Status != "SUCCESS" && result.Status != "FAILURE" {
|
||||
//log.Printf("[DEBUG][%s] Found subflow result that is not handled. Waiting for results", workflowExecution.ExecutionId)
|
||||
|
||||
|
||||
subflowId = result.Action.ID
|
||||
found = true
|
||||
found = true
|
||||
break
|
||||
}
|
||||
}
|
||||
@@ -239,21 +236,20 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
|
||||
if len(subflowId) > 0 {
|
||||
// Under rerun period timeout
|
||||
timeComparison := 120
|
||||
timeComparison := 120
|
||||
log.Printf("[DEBUG][%s] Starting polling for %d seconds to see if new subflow updates are found on the backend that are not handled. Subflow ID: %s", workflowExecution.ExecutionId, timeComparison, subflowId)
|
||||
timestart := time.Now()
|
||||
streamResultUrl := fmt.Sprintf("%s/api/v1/streams/results", baseUrl)
|
||||
for {
|
||||
err = handleSubflowPoller(ctx, workflowExecution, streamResultUrl, subflowId)
|
||||
err = handleSubflowPoller(ctx, workflowExecution, streamResultUrl, subflowId)
|
||||
if err == nil {
|
||||
log.Printf("[DEBUG] Subflow is finished and we are breaking the thingy")
|
||||
|
||||
|
||||
if os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" && workflowExecution.ExecutionSource != "default" {
|
||||
log.Printf("[DEBUG] Force shutdown of worker due to optimized run with webserver. Expecting reruns to take care of this")
|
||||
os.Exit(0)
|
||||
}
|
||||
|
||||
|
||||
break
|
||||
}
|
||||
|
||||
@@ -275,7 +271,6 @@ func setWorkflowExecution(ctx context.Context, workflowExecution shuffle.Workflo
|
||||
return nil
|
||||
}
|
||||
|
||||
|
||||
// removes every container except itself (worker)
|
||||
func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason string, handleResultSend bool) {
|
||||
log.Printf("[DEBUG][%s] Shutdown (%s) started with reason %#v. Result amount: %d. ResultsSent: %d, Send result: %#v, Parent: %#v", workflowExecution.ExecutionId, workflowExecution.Status, reason, len(workflowExecution.Results), requestsSent, handleResultSend, workflowExecution.ExecutionParent)
|
||||
@@ -315,7 +310,7 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
|
||||
}
|
||||
*/
|
||||
} else {
|
||||
|
||||
|
||||
}
|
||||
|
||||
if len(reason) > 0 && len(nodeId) > 0 {
|
||||
@@ -504,7 +499,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
||||
|
||||
if !strings.Contains(param.Value, "shuffle-backend") {
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// Automatic replacement as this is default
|
||||
if len(os.Getenv("BASE_URL")) > 0 {
|
||||
@@ -519,7 +514,6 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// Max 10% CPU every second
|
||||
//CPUShares: 128,
|
||||
//CPUQuota: 10000,
|
||||
@@ -546,7 +540,7 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
||||
|
||||
// Get environment for certificates
|
||||
volumeBinds := []string{}
|
||||
volumeBindString:= os.Getenv("SHUFFLE_VOLUME_BINDS")
|
||||
volumeBindString := os.Getenv("SHUFFLE_VOLUME_BINDS")
|
||||
if len(volumeBindString) > 0 {
|
||||
volumeBindSplit := strings.Split(volumeBindString, ",")
|
||||
for _, volumeBind := range volumeBindSplit {
|
||||
@@ -587,7 +581,6 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
||||
Env: env,
|
||||
}
|
||||
|
||||
|
||||
// Checking as late as possible, just in case.
|
||||
newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
|
||||
_, err := shuffle.GetCache(ctx, newExecId)
|
||||
@@ -867,7 +860,7 @@ func askOtherWorkersToDownloadImage(image string) {
|
||||
// Check environment SHUFFLE_AUTO_IMAGE_DOWNLOAD
|
||||
if os.Getenv("SHUFFLE_AUTO_IMAGE_DOWNLOAD") == "false" {
|
||||
log.Printf("[DEBUG] SHUFFLE_AUTO_IMAGE_DOWNLOAD is false. NOT distributing images %s", image)
|
||||
return
|
||||
return
|
||||
}
|
||||
|
||||
if shuffle.ArrayContains(imagesDistributed, image) {
|
||||
@@ -900,7 +893,7 @@ func askOtherWorkersToDownloadImage(image string) {
|
||||
req, err := http.NewRequest(
|
||||
"POST",
|
||||
url,
|
||||
bytes.NewBuffer(imageJSON),
|
||||
bytes.NewBuffer(imageJSON),
|
||||
)
|
||||
|
||||
if err != nil {
|
||||
@@ -940,7 +933,6 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
startAction, extra, children, parents, visited, executed, nextActions, environments := shuffle.GetExecutionVariables(ctx, workflowExecution.ExecutionId)
|
||||
|
||||
dockercli, err := dockerclient.NewEnvClient()
|
||||
@@ -1008,7 +1000,7 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
||||
|
||||
// marshal action and put it in there rofl
|
||||
//log.Printf("[INFO][%s] Time to execute %s (%s) with app %s:%s, function %s, env %s with %d parameters.", workflowExecution.ExecutionId, action.ID, action.Label, action.AppName, action.AppVersion, action.Name, action.Environment, len(action.Parameters))
|
||||
|
||||
|
||||
log.Printf("[DEBUG][%s] Action: Send, Label: '%s', Action: '%s', Run status: %s, Extra=", workflowExecution.ExecutionId, action.Label, action.AppName, workflowExecution.Status)
|
||||
|
||||
actionData, err := json.Marshal(action)
|
||||
@@ -1094,10 +1086,9 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
||||
}
|
||||
|
||||
if len(os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")) > 0 {
|
||||
env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
|
||||
env = append(env, fmt.Sprintf("SHUFFLE_APP_SDK_TIMEOUT=%s", os.Getenv("SHUFFLE_APP_SDK_TIMEOUT")))
|
||||
}
|
||||
|
||||
|
||||
// Fixes issue:
|
||||
// standard_go init_linux.go:185: exec user process caused "argument list too long"
|
||||
// https://devblogs.microsoft.com/oldnewthing/20100203-00/?p=15083
|
||||
@@ -1125,8 +1116,6 @@ func handleExecutionResult(workflowExecution shuffle.WorkflowExecution) {
|
||||
fmt.Sprintf("%s:%s_%s", baseimagename, parsedAppname, action.AppVersion),
|
||||
}
|
||||
|
||||
|
||||
|
||||
// If cleanup is set, it should run for efficiency
|
||||
pullOptions := types.ImagePullOptions{}
|
||||
if cleanupEnv == "true" {
|
||||
@@ -1629,7 +1618,6 @@ func handleSubflowPoller(ctx context.Context, workflowExecution shuffle.Workflow
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
if workflowExecution.Status == "WAITING" && workflowExecution.ExecutionSource != "default" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm" {
|
||||
log.Printf("[INFO][%s] Workflow execution is waiting. Exiting worker, as backend will restart it.", workflowExecution.ExecutionId)
|
||||
shutdown(workflowExecution, "", "", true)
|
||||
@@ -2062,10 +2050,9 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": true, "reason": "Execution is not executing, but %s"}`, workflowExecution.Status)))
|
||||
}
|
||||
|
||||
|
||||
log.Printf("[DEBUG][%s] Shutting down (35)", workflowExecution.ExecutionId)
|
||||
|
||||
// Force sending result
|
||||
// Force sending result
|
||||
shutdownData, err := json.Marshal(workflowExecution)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR][%s] Failed marshalling execution (35): %s", workflowExecution.ExecutionId, err)
|
||||
@@ -2141,11 +2128,11 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
||||
attempts += 1
|
||||
log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(parsedValue.Results), resultLength)
|
||||
/*
|
||||
if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) {
|
||||
log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions))
|
||||
runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp)
|
||||
return
|
||||
}
|
||||
if len(workflowExecution.Results) <= len(workflowExecution.Workflow.Actions) {
|
||||
log.Printf("[DEBUG][%s] Rerunning transaction as results has changed. %d vs %d", workflowExecution.ExecutionId, len(workflowExecution.Results), len(workflowExecution.Workflow.Actions))
|
||||
runWorkflowExecutionTransaction(ctx, attempts, workflowExecutionId, actionResult, resp)
|
||||
return
|
||||
}
|
||||
*/
|
||||
}
|
||||
}
|
||||
@@ -2178,8 +2165,6 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
||||
}
|
||||
|
||||
func sendSelfRequest(actionResult shuffle.ActionResult) {
|
||||
|
||||
|
||||
data, err := json.Marshal(actionResult)
|
||||
if err != nil {
|
||||
log.Printf("[ERROR][%s] Shutting down (24): Failed to unmarshal data for backend: %s", actionResult.ExecutionId, err)
|
||||
@@ -2238,10 +2223,10 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
||||
|
||||
// Basically to reduce backend strain
|
||||
/*
|
||||
if shuffle.ArrayContains(finishedExecutions, workflowExecution.ExecutionId) {
|
||||
log.Printf("[INFO][%s] NOT sending backend info since it's already been sent before.", workflowExecution.ExecutionId)
|
||||
return
|
||||
}
|
||||
if shuffle.ArrayContains(finishedExecutions, workflowExecution.ExecutionId) {
|
||||
log.Printf("[INFO][%s] NOT sending backend info since it's already been sent before.", workflowExecution.ExecutionId)
|
||||
return
|
||||
}
|
||||
*/
|
||||
|
||||
// Take it down again
|
||||
@@ -2252,6 +2237,7 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
||||
}
|
||||
|
||||
finishedExecutions = append(finishedExecutions, workflowExecution.ExecutionId)
|
||||
|
||||
*/
|
||||
|
||||
streamUrl := fmt.Sprintf("%s/api/v1/streams", baseUrl)
|
||||
@@ -2308,6 +2294,7 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
||||
|
||||
if workflowExecution.Status == "FINISHED" || workflowExecution.Status == "ABORTED" || (len(environments) == 1 && requestsSent == 0 && len(workflowExecution.Results) >= 1 && os.Getenv("SHUFFLE_SWARM_CONFIG") != "run" && os.Getenv("SHUFFLE_SWARM_CONFIG") != "swarm") || (len(workflowExecution.Results) >= len(workflowExecution.Workflow.Actions)+extra && len(workflowExecution.Workflow.Actions) > 0) {
|
||||
|
||||
|
||||
if workflowExecution.Status == "FINISHED" {
|
||||
for _, result := range workflowExecution.Results {
|
||||
if result.Status == "EXECUTING" || result.Status == "WAITING" {
|
||||
@@ -2316,8 +2303,7 @@ func sendResult(workflowExecution shuffle.WorkflowExecution, data []byte) {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
log.Printf("[DEBUG][%s] Should send full result to %s", workflowExecution.ExecutionId, baseUrl)
|
||||
|
||||
//data = fmt.Sprintf(`{"execution_id": "%s", "authorization": "%s"}`, executionId, authorization)
|
||||
@@ -2401,7 +2387,7 @@ func handleGetStreamResults(resp http.ResponseWriter, request *http.Request) {
|
||||
// GetLocalIP returns the non loopback local IP of the host
|
||||
func getLocalIP() string {
|
||||
|
||||
|
||||
|
||||
addrs, err := net.InterfaceAddrs()
|
||||
if err != nil {
|
||||
return ""
|
||||
@@ -2446,8 +2432,7 @@ func webserverSetup(workflowExecution shuffle.WorkflowExecution) net.Listener {
|
||||
}
|
||||
|
||||
log.Printf("[DEBUG] OLD HOSTNAME: %s", appCallbackUrl)
|
||||
|
||||
|
||||
|
||||
port := listener.Addr().(*net.TCPAddr).Port
|
||||
// Set the port environment variable
|
||||
os.Setenv("WORKER_PORT", fmt.Sprintf("%d", port))
|
||||
@@ -2605,7 +2590,6 @@ func downloadDockerImageBackend(client *http.Client, imageName string) error {
|
||||
*/
|
||||
}
|
||||
|
||||
|
||||
// Runs data discovery
|
||||
|
||||
func sendAppRequest(ctx context.Context, incomingUrl, appName string, port int, action *shuffle.Action, workflowExecution *shuffle.WorkflowExecution) error {
|
||||
@@ -2869,7 +2853,7 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe
|
||||
if newresp.StatusCode != 200 {
|
||||
log.Printf("[ERROR] %sStatusCode (1): %d", string(body), newresp.StatusCode)
|
||||
time.Sleep(time.Duration(sleepTime) * time.Second)
|
||||
return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode) )
|
||||
return environments, errors.New(fmt.Sprintf("Bad status code: %d", newresp.StatusCode))
|
||||
}
|
||||
|
||||
err = json.Unmarshal(body, &workflowExecution)
|
||||
@@ -2952,7 +2936,6 @@ func getStreamResultsWrapper(client *http.Client, req *http.Request, workflowExe
|
||||
|
||||
// Set environment variable
|
||||
|
||||
|
||||
//log.Printf("Before wait")
|
||||
//wg := sync.WaitGroup{}
|
||||
//wg.Add(1)
|
||||
@@ -3024,7 +3007,6 @@ func main() {
|
||||
swarmConfig := os.Getenv("SHUFFLE_SWARM_CONFIG")
|
||||
log.Printf("[INFO] Running with timezone %s and swarm config %#v", timezone, swarmConfig)
|
||||
|
||||
|
||||
authorization := ""
|
||||
executionId := ""
|
||||
|
||||
@@ -3329,7 +3311,6 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) {
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
for _, img := range images {
|
||||
for _, tag := range img.RepoTags {
|
||||
splitTag := strings.Split(tag, ":")
|
||||
@@ -3342,7 +3323,7 @@ func handleDownloadImage(resp http.ResponseWriter, request *http.Request) {
|
||||
possibleNames = append(possibleNames, fmt.Sprintf("frikky/shuffle:%s", baseTag))
|
||||
possibleNames = append(possibleNames, fmt.Sprintf("registry.hub.docker.com/frikky/shuffle:%s", baseTag))
|
||||
|
||||
if (arrayContains(possibleNames, image.Image)) {
|
||||
if arrayContains(possibleNames, image.Image) {
|
||||
log.Printf("[DEBUG] Image %s already downloaded that has been requested to download", image.Image)
|
||||
resp.WriteHeader(200)
|
||||
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "image already present"}`)))
|
||||
@@ -3367,7 +3348,6 @@ func runWebserver(listener net.Listener) {
|
||||
r.HandleFunc("/api/v1/run", handleRunExecution).Methods("POST", "OPTIONS")
|
||||
r.HandleFunc("/api/v1/download", handleDownloadImage).Methods("POST", "OPTIONS")
|
||||
|
||||
|
||||
if strings.ToLower(os.Getenv("SHUFFLE_DEBUG_MEMORY")) == "true" {
|
||||
r.HandleFunc("/debug/pprof/", pprof.Index)
|
||||
r.HandleFunc("/debug/pprof/heap", pprof.Handler("heap").ServeHTTP)
|
||||
|
||||
Reference in New Issue
Block a user