worker changes/fixes
This commit is contained in:
@@ -53,7 +53,7 @@ data:
|
|||||||
SSO_REDIRECT_URL: ""
|
SSO_REDIRECT_URL: ""
|
||||||
TZ: "Europe/Amsterdam \t\t\t\t\t"
|
TZ: "Europe/Amsterdam \t\t\t\t\t"
|
||||||
IS_KUBERNETES: "true"
|
IS_KUBERNETES: "true"
|
||||||
REGISTRY_URL: "shuffle-registry:5000"
|
REGISTRY_URL: "172.17.14.71:5000"
|
||||||
#REGISTRY_AUTH: ""
|
#REGISTRY_AUTH: ""
|
||||||
kind: ConfigMap
|
kind: ConfigMap
|
||||||
metadata:
|
metadata:
|
||||||
|
|||||||
@@ -32,6 +32,8 @@ FROM alpine:3.15.0
|
|||||||
ENV SHUFFLE_BASE_IMAGE_REGISTRY=docker.io
|
ENV SHUFFLE_BASE_IMAGE_REGISTRY=docker.io
|
||||||
ENV SHUFFLE_BASE_IMAGE_NAME=frikky/shuffle
|
ENV SHUFFLE_BASE_IMAGE_NAME=frikky/shuffle
|
||||||
ENV SHUFFLE_BASE_IMAGE_TAG_SUFFIX=0.8.70
|
ENV SHUFFLE_BASE_IMAGE_TAG_SUFFIX=0.8.70
|
||||||
|
|
||||||
|
#for k8s
|
||||||
ENV SHUFFLE_OPENSEARCH_URL=https://opensearch:9200
|
ENV SHUFFLE_OPENSEARCH_URL=https://opensearch:9200
|
||||||
ENV SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY=true
|
ENV SHUFFLE_OPENSEARCH_SKIPSSL_VERIFY=true
|
||||||
|
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ import (
|
|||||||
"github.com/docker/docker/api/types"
|
"github.com/docker/docker/api/types"
|
||||||
"github.com/docker/docker/api/types/container"
|
"github.com/docker/docker/api/types/container"
|
||||||
//"github.com/docker/docker/api/types/filters"
|
//"github.com/docker/docker/api/types/filters"
|
||||||
// "github.com/docker/docker/api/types/mount"
|
"github.com/docker/docker/api/types/mount"
|
||||||
dockerclient "github.com/docker/docker/client"
|
dockerclient "github.com/docker/docker/client"
|
||||||
//"github.com/go-git/go-billy/v5/memfs"
|
//"github.com/go-git/go-billy/v5/memfs"
|
||||||
|
|
||||||
@@ -112,18 +112,6 @@ func shutdown(workflowExecution shuffle.WorkflowExecution, nodeId string, reason
|
|||||||
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)
|
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)
|
||||||
//reason := "Error in execution"
|
//reason := "Error in execution"
|
||||||
|
|
||||||
// if os.Getenv("IS_KUBERNETES") == "true" {
|
|
||||||
// log.Printf("[DEBUG][%s] Running in Kubernetes shutting down")
|
|
||||||
// workerPod := fmt.Sprintf("worker-%s", workflowExecution.ExecutionId)
|
|
||||||
// namespace := "shuffle"
|
|
||||||
// clientset, err := getKubernetesClient()
|
|
||||||
// if err != nil {
|
|
||||||
// log.Printf("[ERROR] Error getting kubernetes client: %s", err)
|
|
||||||
// return
|
|
||||||
// }
|
|
||||||
|
|
||||||
// } else {
|
|
||||||
|
|
||||||
sleepDuration := 1
|
sleepDuration := 1
|
||||||
if handleResultSend && requestsSent < 2 {
|
if handleResultSend && requestsSent < 2 {
|
||||||
shutdownData, err := json.Marshal(workflowExecution)
|
shutdownData, err := json.Marshal(workflowExecution)
|
||||||
@@ -237,12 +225,10 @@ func getKubernetesClient() (*kubernetes.Clientset, error) {
|
|||||||
|
|
||||||
// Deploys the internal worker whenever something happens
|
// Deploys the internal worker whenever something happens
|
||||||
func deployApp(cli *dockerclient.Client, image string, identifier string, env []string, workflowExecution shuffle.WorkflowExecution, action shuffle.Action) error {
|
func deployApp(cli *dockerclient.Client, image string, identifier string, env []string, workflowExecution shuffle.WorkflowExecution, action shuffle.Action) error {
|
||||||
log.Printf("HELLO FROM DEPLOYAPP")
|
log.Printf("################################### new call to deployApp ###################################")
|
||||||
log.Printf("image: %s", image)
|
log.Printf("image: %s", image)
|
||||||
log.Printf("identifier: %s", identifier)
|
log.Printf("identifier: %s", identifier)
|
||||||
|
// log.Printf("execution: %+v", workflowExecution)
|
||||||
log.Printf("IS_KUBERNETES: %s", os.Getenv("IS_KUBERNETES"))
|
|
||||||
log.Printf("REGISTRY_NAME: %s", os.Getenv("REGISTRY_NAME"))
|
|
||||||
|
|
||||||
if os.Getenv("IS_KUBERNETES") == "true" {
|
if os.Getenv("IS_KUBERNETES") == "true" {
|
||||||
|
|
||||||
@@ -268,6 +254,26 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|||||||
value := strSplit[0]
|
value := strSplit[0]
|
||||||
value = strings.ReplaceAll(value, "_", "-")
|
value = strings.ReplaceAll(value, "_", "-")
|
||||||
|
|
||||||
|
// checking if app is generated or not
|
||||||
|
appDetails := strings.Split(image, ":")[1]
|
||||||
|
appDetailsSplit := strings.Split(appDetails, "_")
|
||||||
|
appName := appDetailsSplit[0]
|
||||||
|
appVersion := appDetailsSplit[1]
|
||||||
|
|
||||||
|
for _, app := range workflowExecution.Workflow.Actions {
|
||||||
|
log.Printf("[DEBUG] App: %s, Version: %s", appName, appVersion)
|
||||||
|
log.Printf("[DEBUG] Checking app %s with version %s", app.AppName, app.AppVersion)
|
||||||
|
if app.AppName == appName && app.AppVersion == appVersion {
|
||||||
|
if app.Generated == true {
|
||||||
|
log.Printf("[DEBUG] Generated app, setting local registry")
|
||||||
|
image = fmt.Sprintf("%s/%s", registryName, image)
|
||||||
|
break
|
||||||
|
} else {
|
||||||
|
log.Printf("[DEBUG] Not generated app, setting shuffle registry")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
//fix naming convention
|
//fix naming convention
|
||||||
podUuid := uuid.NewV4().String()
|
podUuid := uuid.NewV4().String()
|
||||||
podName := fmt.Sprintf("%s-%s", value, podUuid)
|
podName := fmt.Sprintf("%s-%s", value, podUuid)
|
||||||
@@ -297,11 +303,112 @@ func deployApp(cli *dockerclient.Client, image string, identifier string, env []
|
|||||||
fmt.Fprintf(os.Stderr, "Error creating pod: %v\n", err)
|
fmt.Fprintf(os.Stderr, "Error creating pod: %v\n", err)
|
||||||
os.Exit(1)
|
os.Exit(1)
|
||||||
}
|
}
|
||||||
fmt.Printf("Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace)
|
fmt.Printf("[DEBUG] Created pod %q in namespace %q\n", createdPod.Name, createdPod.Namespace)
|
||||||
} else {
|
} else {
|
||||||
// docker part
|
// form basic hostConfig
|
||||||
|
ctx := context.Background()
|
||||||
|
|
||||||
|
if action.AppName == "shuffle-subflow" {
|
||||||
|
// Automatic replacement of URL
|
||||||
|
for paramIndex, param := range action.Parameters {
|
||||||
|
if param.Name != "backend_url" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
if strings.Contains(param.Value, "shuffle-backend") {
|
||||||
|
// Automatic replacement as this is default
|
||||||
|
action.Parameters[paramIndex].Value = os.Getenv("BASE_URL")
|
||||||
|
log.Printf("[DEBUG][%s] Replaced backend_url with %s", workflowExecution.ExecutionId, os.Getenv("BASE_URL"))
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Max 10% CPU every second
|
||||||
|
//CPUShares: 128,
|
||||||
|
//CPUQuota: 10000,
|
||||||
|
//CPUPeriod: 100000,
|
||||||
|
hostConfig := &container.HostConfig{
|
||||||
|
LogConfig: container.LogConfig{
|
||||||
|
Type: "json-file",
|
||||||
|
Config: map[string]string{
|
||||||
|
"max-size": "10m",
|
||||||
|
},
|
||||||
|
},
|
||||||
|
Resources: container.Resources{},
|
||||||
|
}
|
||||||
|
|
||||||
|
hostConfig.NetworkMode = container.NetworkMode(fmt.Sprintf("container:worker-%s", workflowExecution.ExecutionId))
|
||||||
|
|
||||||
|
// Removing because log extraction should happen first
|
||||||
|
if cleanupEnv == "true" {
|
||||||
|
hostConfig.AutoRemove = true
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIXME: Add proper foldermounts here
|
||||||
|
//log.Printf("\n\nPRE FOLDERMOUNT\n\n")
|
||||||
|
//volumeBinds := []string{"/tmp/shuffle-mount:/rules"}
|
||||||
|
//volumeBinds := []string{"/tmp/shuffle-mount:/rules"}
|
||||||
|
volumeBinds := []string{}
|
||||||
|
if len(volumeBinds) > 0 {
|
||||||
|
log.Printf("[DEBUG] Setting up binds for container!")
|
||||||
|
hostConfig.Binds = volumeBinds
|
||||||
|
hostConfig.Mounts = []mount.Mount{}
|
||||||
|
for _, bind := range volumeBinds {
|
||||||
|
if !strings.Contains(bind, ":") || strings.Contains(bind, "..") || strings.HasPrefix(bind, "~") {
|
||||||
|
log.Printf("[WARNING] Bind %s is invalid.", bind)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
log.Printf("[DEBUG] Appending bind %s", bind)
|
||||||
|
bindSplit := strings.Split(bind, ":")
|
||||||
|
sourceFolder := bindSplit[0]
|
||||||
|
destinationFolder := bindSplit[0]
|
||||||
|
hostConfig.Mounts = append(hostConfig.Mounts, mount.Mount{
|
||||||
|
Type: mount.TypeBind,
|
||||||
|
Source: sourceFolder,
|
||||||
|
Target: destinationFolder,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
//log.Printf("[WARNING] Not mounting folders")
|
||||||
|
}
|
||||||
|
|
||||||
|
config := &container.Config{
|
||||||
|
Image: image,
|
||||||
|
Env: env,
|
||||||
|
}
|
||||||
|
|
||||||
|
// Checking as late as possible, just in case.
|
||||||
|
newExecId := fmt.Sprintf("%s_%s", workflowExecution.ExecutionId, action.ID)
|
||||||
|
_, err := shuffle.GetCache(ctx, newExecId)
|
||||||
|
if err == nil {
|
||||||
|
log.Printf("\n\n[DEBUG] Result for %s already found - returning\n\n", newExecId)
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
cacheData := []byte("1")
|
||||||
|
err = shuffle.SetCache(ctx, newExecId, cacheData, 30)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[WARNING] Failed setting cache for action %s: %s", newExecId, err)
|
||||||
|
} else {
|
||||||
|
log.Printf("[DEBUG] Adding %s to cache. Name: %s", newExecId, action.Name)
|
||||||
|
}
|
||||||
|
|
||||||
|
if action.ExecutionDelay > 0 {
|
||||||
|
log.Printf("[DEBUG] Running app %s in docker with delay of %d", action.Name, action.ExecutionDelay)
|
||||||
|
waitTime := time.Duration(action.ExecutionDelay) * time.Second
|
||||||
|
|
||||||
|
time.AfterFunc(waitTime, func() {
|
||||||
|
DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
log.Printf("[DEBUG] Running app %s in docker NORMALLY as there is no delay set with identifier %s", action.Name, identifier)
|
||||||
|
returnvalue := DeployContainer(ctx, cli, config, hostConfig, identifier, workflowExecution, newExecId)
|
||||||
|
log.Printf("[DEBUG] Normal deploy ret: %s", returnvalue)
|
||||||
|
return returnvalue
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -1168,7 +1275,17 @@ func handleDefaultExecution(client *http.Client, req *http.Request, workflowExec
|
|||||||
if workflowExecution.Status != "EXECUTING" {
|
if workflowExecution.Status != "EXECUTING" {
|
||||||
log.Printf("[WARNING][%s] Exiting as worker execution has status %s!", workflowExecution.ExecutionId, workflowExecution.Status)
|
log.Printf("[WARNING][%s] Exiting as worker execution has status %s!", workflowExecution.ExecutionId, workflowExecution.Status)
|
||||||
log.Printf("[DEBUG] Shutting down (21)")
|
log.Printf("[DEBUG] Shutting down (21)")
|
||||||
shutdown(workflowExecution, "", "", true)
|
if os.Getenv("IS_KUBERNETES") == "true" {
|
||||||
|
// log.Printf("workflow execution: %#v", workflowExecution)
|
||||||
|
clientset, err := getKubernetesClient()
|
||||||
|
if err != nil {
|
||||||
|
fmt.Println("[ERROR]Error getting kubernetes client:", err)
|
||||||
|
os.Exit(1)
|
||||||
|
}
|
||||||
|
cleanupExecution(clientset, workflowExecution, "shuffle")
|
||||||
|
} else {
|
||||||
|
shutdown(workflowExecution, "", "", true)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
setWorkflowExecution(ctx, workflowExecution, false)
|
setWorkflowExecution(ctx, workflowExecution, false)
|
||||||
|
|||||||
Reference in New Issue
Block a user