Fixed merge conflicts
@@ -534,7 +534,7 @@ func getDockerImage(resp http.ResponseWriter, request *http.Request) {
|
|||||||
if len(img.ID) == 0 {
|
if len(img.ID) == 0 {
|
||||||
if len(img2.ID) == 0 {
|
if len(img2.ID) == 0 {
|
||||||
workflowapps, err := shuffle.GetAllWorkflowApps(ctx, 0, 0)
|
workflowapps, err := shuffle.GetAllWorkflowApps(ctx, 0, 0)
|
||||||
//log.Printf("[INFO] Getting workflowapps for a rebuild. Got %d with err %#v", len(workflowapps), err)
|
log.Printf("[INFO] Getting workflowapps for a rebuild. Got %d with err %#v", len(workflowapps), err)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
imageName := ""
|
imageName := ""
|
||||||
imageVersion := ""
|
imageVersion := ""
|
||||||
|
|||||||
@@ -49,7 +49,6 @@ require (
|
|||||||
github.com/docker/go-connections v0.4.0 // indirect
|
github.com/docker/go-connections v0.4.0 // indirect
|
||||||
github.com/docker/go-units v0.5.0 // indirect
|
github.com/docker/go-units v0.5.0 // indirect
|
||||||
github.com/emirpasic/gods v1.18.1 // indirect
|
github.com/emirpasic/gods v1.18.1 // indirect
|
||||||
github.com/frikky/go-elasticsearch/v8 v8.13.1 // indirect
|
|
||||||
github.com/go-git/gcfg v1.5.1-0.20230307220236-3a3c6141e376 // indirect
|
github.com/go-git/gcfg v1.5.1-0.20230307220236-3a3c6141e376 // indirect
|
||||||
github.com/go-openapi/jsonpointer v0.19.5 // indirect
|
github.com/go-openapi/jsonpointer v0.19.5 // indirect
|
||||||
github.com/go-openapi/swag v0.19.5 // indirect
|
github.com/go-openapi/swag v0.19.5 // indirect
|
||||||
@@ -102,4 +101,5 @@ require (
|
|||||||
google.golang.org/protobuf v1.30.0 // indirect
|
google.golang.org/protobuf v1.30.0 // indirect
|
||||||
gopkg.in/warnings.v0 v0.1.2 // indirect
|
gopkg.in/warnings.v0 v0.1.2 // indirect
|
||||||
gopkg.in/yaml.v2 v2.4.0 // indirect
|
gopkg.in/yaml.v2 v2.4.0 // indirect
|
||||||
|
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -5997,6 +5997,7 @@ func initHandlers() {
|
|||||||
// New for recommendations in Shuffle
|
// New for recommendations in Shuffle
|
||||||
r.HandleFunc("/api/v1/recommendations/get_actions", shuffle.HandleActionRecommendation).Methods("POST", "OPTIONS")
|
r.HandleFunc("/api/v1/recommendations/get_actions", shuffle.HandleActionRecommendation).Methods("POST", "OPTIONS")
|
||||||
r.HandleFunc("/api/v1/recommendations/modify", shuffle.HandleRecommendationAction).Methods("POST", "OPTIONS")
|
r.HandleFunc("/api/v1/recommendations/modify", shuffle.HandleRecommendationAction).Methods("POST", "OPTIONS")
|
||||||
|
r.HandleFunc("/api/v1/workflows/{key}/revisions", shuffle.GetWorkflowRevisions).Methods("GET", "OPTIONS")
|
||||||
|
|
||||||
// Triggers
|
// Triggers
|
||||||
r.HandleFunc("/api/v1/hooks/new", shuffle.HandleNewHook).Methods("POST", "OPTIONS")
|
r.HandleFunc("/api/v1/hooks/new", shuffle.HandleNewHook).Methods("POST", "OPTIONS")
|
||||||
|
|||||||
@@ -636,6 +636,7 @@ func handleWorkflowQueue(resp http.ResponseWriter, request *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
//log.Printf("Actionresult unmarshal: %s", string(body))
|
//log.Printf("Actionresult unmarshal: %s", string(body))
|
||||||
|
log.Printf("[DEBUG] Got workflow result from %s of length %d", request.RemoteAddr, len(body))
|
||||||
ctx := context.Background()
|
ctx := context.Background()
|
||||||
err = shuffle.ValidateNewWorkerExecution(ctx, body)
|
err = shuffle.ValidateNewWorkerExecution(ctx, body)
|
||||||
if err == nil {
|
if err == nil {
|
||||||
@@ -774,7 +775,6 @@ func runWorkflowExecutionTransaction(ctx context.Context, attempts int64, workfl
|
|||||||
}
|
}
|
||||||
|
|
||||||
//log.Printf("BASE LENGTH: %d", len(workflowExecution.Results))
|
//log.Printf("BASE LENGTH: %d", len(workflowExecution.Results))
|
||||||
|
|
||||||
workflowExecution, dbSave, err := shuffle.ParsedExecutionResult(ctx, *workflowExecution, actionResult, false, 0)
|
workflowExecution, dbSave, err := shuffle.ParsedExecutionResult(ctx, *workflowExecution, actionResult, false, 0)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
b, suberr := json.Marshal(actionResult)
|
b, suberr := json.Marshal(actionResult)
|
||||||
@@ -1118,6 +1118,678 @@ func handleExecution(id string, workflow shuffle.Workflow, request *http.Request
|
|||||||
err = imageCheckBuilder(execInfo.ImageNames)
|
err = imageCheckBuilder(execInfo.ImageNames)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
log.Printf("[ERROR] Failed building the required images from %#v: %s", execInfo.ImageNames, err)
|
log.Printf("[ERROR] Failed building the required images from %#v: %s", execInfo.ImageNames, err)
|
||||||
|
return shuffle.WorkflowExecution{}, "Failed unmarshal during execution", err
|
||||||
|
}
|
||||||
|
|
||||||
|
makeNew := true
|
||||||
|
start, startok := request.URL.Query()["start"]
|
||||||
|
if request.Method == "POST" {
|
||||||
|
body, err := ioutil.ReadAll(request.Body)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] Failed request POST read: %s", err)
|
||||||
|
return shuffle.WorkflowExecution{}, "Failed getting body", err
|
||||||
|
}
|
||||||
|
|
||||||
|
// This one doesn't really matter.
|
||||||
|
log.Printf("[INFO] Running POST execution with body of length %d for workflow %s", len(string(body)), workflowExecution.Workflow.ID)
|
||||||
|
|
||||||
|
if len(body) >= 4 {
|
||||||
|
if body[0] == 34 && body[len(body)-1] == 34 {
|
||||||
|
body = body[1 : len(body)-1]
|
||||||
|
}
|
||||||
|
if body[0] == 34 && body[len(body)-1] == 34 {
|
||||||
|
body = body[1 : len(body)-1]
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
sourceAuth, sourceAuthOk := request.URL.Query()["source_auth"]
|
||||||
|
if sourceAuthOk {
|
||||||
|
//log.Printf("\n\n\nSETTING SOURCE WORKFLOW AUTH TO %s!!!\n\n\n", sourceAuth[0])
|
||||||
|
workflowExecution.ExecutionSourceAuth = sourceAuth[0]
|
||||||
|
} else {
|
||||||
|
//log.Printf("Did NOT get source workflow")
|
||||||
|
}
|
||||||
|
|
||||||
|
sourceNode, sourceNodeOk := request.URL.Query()["source_node"]
|
||||||
|
if sourceNodeOk {
|
||||||
|
//log.Printf("\n\n\nSETTING SOURCE WORKFLOW NODE TO %s!!!\n\n\n", sourceNode[0])
|
||||||
|
workflowExecution.ExecutionSourceNode = sourceNode[0]
|
||||||
|
} else {
|
||||||
|
//log.Printf("Did NOT get source workflow")
|
||||||
|
}
|
||||||
|
|
||||||
|
//workflowExecution.ExecutionSource = "default"
|
||||||
|
sourceWorkflow, sourceWorkflowOk := request.URL.Query()["source_workflow"]
|
||||||
|
if sourceWorkflowOk {
|
||||||
|
//log.Printf("Got source workflow %s", sourceWorkflow)
|
||||||
|
workflowExecution.ExecutionSource = sourceWorkflow[0]
|
||||||
|
} else {
|
||||||
|
//log.Printf("Did NOT get source workflow")
|
||||||
|
}
|
||||||
|
|
||||||
|
sourceExecution, sourceExecutionOk := request.URL.Query()["source_execution"]
|
||||||
|
if sourceExecutionOk {
|
||||||
|
//log.Printf("[INFO] Got source execution%s", sourceExecution)
|
||||||
|
workflowExecution.ExecutionParent = sourceExecution[0]
|
||||||
|
} else {
|
||||||
|
//log.Printf("Did NOT get source execution")
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(string(body)) < 50 {
|
||||||
|
//log.Println(body)
|
||||||
|
// String in string
|
||||||
|
//log.Println(body)
|
||||||
|
|
||||||
|
//if string(body)[0] == "\"" && string(body)[string(body)
|
||||||
|
log.Printf("[DEBUG] Body: %s", string(body))
|
||||||
|
}
|
||||||
|
|
||||||
|
var execution shuffle.ExecutionRequest
|
||||||
|
err = json.Unmarshal(body, &execution)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[WARNING] Failed execution POST unmarshaling - continuing anyway: %s", err)
|
||||||
|
//return shuffle.WorkflowExecution{}, "", err
|
||||||
|
}
|
||||||
|
|
||||||
|
if execution.Start == "" && len(body) > 0 {
|
||||||
|
execution.ExecutionArgument = string(body)
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIXME - this should have "execution_argument" from executeWorkflow frontend
|
||||||
|
//log.Printf("EXEC: %#v", execution)
|
||||||
|
if len(execution.ExecutionArgument) > 0 {
|
||||||
|
workflowExecution.ExecutionArgument = execution.ExecutionArgument
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(execution.ExecutionSource) > 0 {
|
||||||
|
workflowExecution.ExecutionSource = execution.ExecutionSource
|
||||||
|
}
|
||||||
|
|
||||||
|
//log.Printf("Execution data: %#v", execution)
|
||||||
|
if len(execution.Start) == 36 && len(workflow.Actions) > 0 {
|
||||||
|
log.Printf("[INFO] Should start execution on node %s", execution.Start)
|
||||||
|
workflowExecution.Start = execution.Start
|
||||||
|
|
||||||
|
found := false
|
||||||
|
for _, action := range workflow.Actions {
|
||||||
|
if action.ID == execution.Start {
|
||||||
|
found = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if !found {
|
||||||
|
log.Printf("[ERROR] Action %s was NOT found! Exiting execution.", execution.Start)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Startnode %s was not found in actions", workflow.Start), errors.New(fmt.Sprintf("Startnode %s was not found in actions", workflow.Start))
|
||||||
|
}
|
||||||
|
} else if len(execution.Start) > 0 {
|
||||||
|
//log.Printf("[INFO] !")
|
||||||
|
//log.Printf("[ERROR] START ACTION %s IS WRONG ID LENGTH %d!", execution.Start, len(execution.Start))
|
||||||
|
//return shuffle.WorkflowExecution{}, fmt.Sprintf("Startnode %s was not found in actions", execution.Start), errors.New(fmt.Sprintf("Startnode %s was not found in actions", execution.Start))
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(execution.ExecutionId) == 36 {
|
||||||
|
workflowExecution.ExecutionId = execution.ExecutionId
|
||||||
|
} else {
|
||||||
|
sessionToken := uuid.NewV4()
|
||||||
|
workflowExecution.ExecutionId = sessionToken.String()
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// Check for parameters of start and ExecutionId
|
||||||
|
// This is mostly used for user input trigger
|
||||||
|
|
||||||
|
answer, answerok := request.URL.Query()["answer"]
|
||||||
|
referenceId, referenceok := request.URL.Query()["reference_execution"]
|
||||||
|
if answerok && referenceok {
|
||||||
|
// If answer is false, reference execution with result
|
||||||
|
log.Printf("[INFO] Answer is OK AND reference is OK!")
|
||||||
|
if answer[0] == "false" {
|
||||||
|
log.Printf("Should update reference and return, no need for further execution!")
|
||||||
|
|
||||||
|
// Get the reference execution
|
||||||
|
oldExecution, err := shuffle.GetWorkflowExecution(ctx, referenceId[0])
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("Failed getting execution (execution) %s: %s", referenceId[0], err)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Failed getting execution ID %s because it doesn't exist.", referenceId[0]), err
|
||||||
|
}
|
||||||
|
|
||||||
|
if oldExecution.Workflow.ID != id {
|
||||||
|
log.Println("Wrong workflowid!")
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Bad ID %s", referenceId), errors.New("Bad ID")
|
||||||
|
}
|
||||||
|
|
||||||
|
newResults := []shuffle.ActionResult{}
|
||||||
|
//log.Printf("%#v", oldExecution.Results)
|
||||||
|
for _, result := range oldExecution.Results {
|
||||||
|
log.Printf("%s - %s", result.Action.ID, start[0])
|
||||||
|
if result.Action.ID == start[0] {
|
||||||
|
note, noteok := request.URL.Query()["note"]
|
||||||
|
if noteok {
|
||||||
|
result.Result = fmt.Sprintf("User note: %s", note[0])
|
||||||
|
} else {
|
||||||
|
result.Result = fmt.Sprintf("User clicked %s", answer[0])
|
||||||
|
}
|
||||||
|
|
||||||
|
// Stopping the whole thing
|
||||||
|
result.CompletedAt = int64(time.Now().Unix())
|
||||||
|
result.Status = "ABORTED"
|
||||||
|
oldExecution.Status = result.Status
|
||||||
|
oldExecution.Result = result.Result
|
||||||
|
oldExecution.LastNode = result.Action.ID
|
||||||
|
}
|
||||||
|
|
||||||
|
newResults = append(newResults, result)
|
||||||
|
}
|
||||||
|
|
||||||
|
oldExecution.Results = newResults
|
||||||
|
err = shuffle.SetWorkflowExecution(ctx, *oldExecution, true)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("Error saving workflow execution actionresult setting: %s", err)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Failed setting workflowexecution actionresult in execution: %s", err), err
|
||||||
|
}
|
||||||
|
|
||||||
|
return shuffle.WorkflowExecution{}, "", nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if referenceok {
|
||||||
|
log.Printf("Handling an old execution continuation!")
|
||||||
|
// Will use the old name, but still continue with NEW ID
|
||||||
|
oldExecution, err := shuffle.GetWorkflowExecution(ctx, referenceId[0])
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("Failed getting execution (execution) %s: %s", referenceId[0], err)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Failed getting execution ID %s because it doesn't exist.", referenceId[0]), err
|
||||||
|
}
|
||||||
|
|
||||||
|
workflowExecution = *oldExecution
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(workflowExecution.ExecutionId) == 0 {
|
||||||
|
sessionToken := uuid.NewV4()
|
||||||
|
workflowExecution.ExecutionId = sessionToken.String()
|
||||||
|
} else {
|
||||||
|
log.Printf("Using the same executionId as before: %s", workflowExecution.ExecutionId)
|
||||||
|
makeNew = false
|
||||||
|
}
|
||||||
|
|
||||||
|
// Don't override workflow defaults
|
||||||
|
}
|
||||||
|
|
||||||
|
if startok {
|
||||||
|
//log.Printf("\n\n[INFO] Setting start to %s based on query!\n\n", start[0])
|
||||||
|
//workflowExecution.Workflow.Start = start[0]
|
||||||
|
workflowExecution.Start = start[0]
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIXME - regex uuid, and check if already exists?
|
||||||
|
if len(workflowExecution.ExecutionId) != 36 {
|
||||||
|
log.Printf("Invalid uuid: %s", workflowExecution.ExecutionId)
|
||||||
|
return shuffle.WorkflowExecution{}, "Invalid uuid", err
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIXME - find owner of workflow
|
||||||
|
// FIXME - get the actual workflow itself and build the request
|
||||||
|
// MAYBE: Don't send the workflow within the pubsub, as this requires more data to be sent
|
||||||
|
// Check if a worker already exists for company, else run one with:
|
||||||
|
// locations, project IDs and subscription names
|
||||||
|
|
||||||
|
// When app is executed:
|
||||||
|
// Should update with status execution (somewhere), which will trigger the next node
|
||||||
|
// IF action.type == internal, we need the internal watcher to be running and executing
|
||||||
|
// This essentially means the WORKER has to be the responsible party for new actions in the INTERNAL landscape
|
||||||
|
// Results are ALWAYS posted back to cloud@execution_id?
|
||||||
|
if makeNew {
|
||||||
|
workflowExecution.Type = "workflow"
|
||||||
|
//workflowExecution.Stream = "tmp"
|
||||||
|
//workflowExecution.WorkflowQueue = "tmp"
|
||||||
|
//workflowExecution.SubscriptionNameNodestream = "testcompany-nodestream"
|
||||||
|
//workflowExecution.Locations = []string{"europe-west2"}
|
||||||
|
workflowExecution.ProjectId = gceProject
|
||||||
|
workflowExecution.WorkflowId = workflow.ID
|
||||||
|
workflowExecution.StartedAt = int64(time.Now().Unix())
|
||||||
|
workflowExecution.CompletedAt = 0
|
||||||
|
workflowExecution.Authorization = uuid.NewV4().String()
|
||||||
|
|
||||||
|
// Status for the entire workflow.
|
||||||
|
workflowExecution.Status = "EXECUTING"
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(workflowExecution.ExecutionSource) == 0 {
|
||||||
|
log.Printf("[INFO] No execution source (trigger) specified. Setting to default")
|
||||||
|
workflowExecution.ExecutionSource = "default"
|
||||||
|
} else {
|
||||||
|
log.Printf("[INFO] Execution source is %s for execution ID %s in workflow %s", workflowExecution.ExecutionSource, workflowExecution.ExecutionId, workflowExecution.Workflow.ID)
|
||||||
|
}
|
||||||
|
|
||||||
|
workflowExecution.ExecutionVariables = workflow.ExecutionVariables
|
||||||
|
if len(workflowExecution.Start) == 0 && len(workflowExecution.Workflow.Start) > 0 {
|
||||||
|
workflowExecution.Start = workflowExecution.Workflow.Start
|
||||||
|
}
|
||||||
|
|
||||||
|
startnodeFound := false
|
||||||
|
newStartnode := ""
|
||||||
|
for _, item := range workflowExecution.Workflow.Actions {
|
||||||
|
if item.ID == workflowExecution.Start {
|
||||||
|
startnodeFound = true
|
||||||
|
}
|
||||||
|
|
||||||
|
if item.IsStartNode {
|
||||||
|
newStartnode = item.ID
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if !startnodeFound {
|
||||||
|
log.Printf("[INFO] Couldn't find startnode %s. Remapping to %#v", workflowExecution.Start, newStartnode)
|
||||||
|
|
||||||
|
if len(newStartnode) > 0 {
|
||||||
|
workflowExecution.Start = newStartnode
|
||||||
|
} else {
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Startnode couldn't be found"), errors.New("Startnode isn't defined in this workflow..")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
childNodes := shuffle.FindChildNodes(workflowExecution, workflowExecution.Start, []string{}, []string{})
|
||||||
|
|
||||||
|
topic := "workflows"
|
||||||
|
startFound := false
|
||||||
|
// FIXME - remove this?
|
||||||
|
newActions := []shuffle.Action{}
|
||||||
|
defaultResults := []shuffle.ActionResult{}
|
||||||
|
|
||||||
|
allAuths := []shuffle.AppAuthenticationStorage{}
|
||||||
|
for _, action := range workflowExecution.Workflow.Actions {
|
||||||
|
//action.LargeImage = ""
|
||||||
|
if action.ID == workflowExecution.Start {
|
||||||
|
startFound = true
|
||||||
|
}
|
||||||
|
//log.Println(action.Environment)
|
||||||
|
|
||||||
|
if action.Environment == "" {
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Environment is not defined for %s", action.Name), errors.New("Environment not defined!")
|
||||||
|
}
|
||||||
|
|
||||||
|
// FIXME: Authentication parameters
|
||||||
|
if len(action.AuthenticationId) > 0 {
|
||||||
|
if len(allAuths) == 0 {
|
||||||
|
allAuths, err = shuffle.GetAllWorkflowAppAuth(ctx, workflow.ExecutingOrg.Id)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("Api authentication failed in get all app auth: %s", err)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Api authentication failed in get all app auth: %s", err), err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
curAuth := shuffle.AppAuthenticationStorage{Id: ""}
|
||||||
|
for _, auth := range allAuths {
|
||||||
|
if auth.Id == action.AuthenticationId {
|
||||||
|
curAuth = auth
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(curAuth.Id) == 0 {
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Auth ID %s doesn't exist", action.AuthenticationId), errors.New(fmt.Sprintf("Auth ID %s doesn't exist", action.AuthenticationId))
|
||||||
|
}
|
||||||
|
|
||||||
|
if curAuth.Encrypted {
|
||||||
|
setField := true
|
||||||
|
newFields := []shuffle.AuthenticationStore{}
|
||||||
|
for _, field := range curAuth.Fields {
|
||||||
|
parsedKey := fmt.Sprintf("%s_%d_%s_%s", curAuth.OrgId, curAuth.Created, curAuth.Label, field.Key)
|
||||||
|
newValue, err := shuffle.HandleKeyDecryption([]byte(field.Value), parsedKey)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[WARNING] Failed decryption for %s: %s", field.Key, err)
|
||||||
|
setField = false
|
||||||
|
break
|
||||||
|
}
|
||||||
|
|
||||||
|
field.Value = string(newValue)
|
||||||
|
newFields = append(newFields, field)
|
||||||
|
}
|
||||||
|
|
||||||
|
if setField {
|
||||||
|
curAuth.Fields = newFields
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
log.Printf("[INFO] AUTH IS NOT ENCRYPTED - attempting encrypting!")
|
||||||
|
err = shuffle.SetWorkflowAppAuthDatastore(ctx, curAuth, curAuth.Id)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[WARNING] Failed running encryption during execution: %s", err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
newParams := []shuffle.WorkflowAppActionParameter{}
|
||||||
|
if strings.ToLower(curAuth.Type) == "oauth2" {
|
||||||
|
log.Printf("[DEBUG] Should replace auth parameters (Oauth2)")
|
||||||
|
|
||||||
|
for _, param := range curAuth.Fields {
|
||||||
|
if param.Key == "expiration" {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
newParams = append(newParams, shuffle.WorkflowAppActionParameter{
|
||||||
|
Name: param.Key,
|
||||||
|
Value: param.Value,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, param := range action.Parameters {
|
||||||
|
//log.Printf("Param: %#v", param)
|
||||||
|
if param.Configuration {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
newParams = append(newParams, param)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
// Rebuild params with the right data. This is to prevent issues on the frontend
|
||||||
|
for _, param := range action.Parameters {
|
||||||
|
|
||||||
|
for _, authparam := range curAuth.Fields {
|
||||||
|
if param.Name == authparam.Key {
|
||||||
|
param.Value = authparam.Value
|
||||||
|
//log.Printf("Name: %s - value: %s", param.Name, param.Value)
|
||||||
|
//log.Printf("Name: %s - value: %s\n", param.Name, param.Value)
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
newParams = append(newParams, param)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
action.Parameters = newParams
|
||||||
|
}
|
||||||
|
|
||||||
|
action.LargeImage = ""
|
||||||
|
if len(action.Label) == 0 {
|
||||||
|
action.Label = action.ID
|
||||||
|
}
|
||||||
|
//log.Printf("LABEL: %s", action.Label)
|
||||||
|
newActions = append(newActions, action)
|
||||||
|
|
||||||
|
// If the node is NOT found, it's supposed to be set to SKIPPED,
|
||||||
|
// as it's not a childnode of the startnode
|
||||||
|
// This is a configuration item for the workflow itself.
|
||||||
|
if len(workflowExecution.Results) > 0 {
|
||||||
|
defaultResults = []shuffle.ActionResult{}
|
||||||
|
for _, result := range workflowExecution.Results {
|
||||||
|
if result.Status == "WAITING" {
|
||||||
|
result.Status = "FINISHED"
|
||||||
|
result.Result = "Continuing"
|
||||||
|
}
|
||||||
|
|
||||||
|
defaultResults = append(defaultResults, result)
|
||||||
|
}
|
||||||
|
} else if len(workflowExecution.Results) == 0 && !workflowExecution.Workflow.Configuration.StartFromTop {
|
||||||
|
found := false
|
||||||
|
for _, nodeId := range childNodes {
|
||||||
|
if nodeId == action.ID {
|
||||||
|
//log.Printf("Found %s", action.ID)
|
||||||
|
found = true
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if !found {
|
||||||
|
if action.ID == workflowExecution.Start {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
//log.Printf("[WARNING] Set %s to SKIPPED as it's NOT a childnode of the startnode.", action.ID)
|
||||||
|
curaction := shuffle.Action{
|
||||||
|
AppName: action.AppName,
|
||||||
|
AppVersion: action.AppVersion,
|
||||||
|
Label: action.Label,
|
||||||
|
Name: action.Name,
|
||||||
|
ID: action.ID,
|
||||||
|
}
|
||||||
|
//action
|
||||||
|
//curaction.Parameters = []
|
||||||
|
defaultResults = append(defaultResults, shuffle.ActionResult{
|
||||||
|
Action: curaction,
|
||||||
|
ExecutionId: workflowExecution.ExecutionId,
|
||||||
|
Authorization: workflowExecution.Authorization,
|
||||||
|
Result: "Skipped because it's not under the startnode",
|
||||||
|
StartedAt: 0,
|
||||||
|
CompletedAt: 0,
|
||||||
|
Status: "SKIPPED",
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
removeTriggers := []string{}
|
||||||
|
for triggerIndex, trigger := range workflowExecution.Workflow.Triggers {
|
||||||
|
//log.Printf("[INFO] ID: %s vs %s", trigger.ID, workflowExecution.Start)
|
||||||
|
if trigger.ID == workflowExecution.Start {
|
||||||
|
if trigger.AppName == "User Input" {
|
||||||
|
startFound = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if trigger.AppName == "User Input" || trigger.AppName == "Shuffle Workflow" {
|
||||||
|
found := false
|
||||||
|
for _, node := range childNodes {
|
||||||
|
if node == trigger.ID {
|
||||||
|
found = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if !found {
|
||||||
|
//log.Printf("SHOULD SET TRIGGER %s TO BE SKIPPED", trigger.ID)
|
||||||
|
|
||||||
|
curaction := shuffle.Action{
|
||||||
|
AppName: "shuffle-subflow",
|
||||||
|
AppVersion: trigger.AppVersion,
|
||||||
|
Label: trigger.Label,
|
||||||
|
Name: trigger.Name,
|
||||||
|
ID: trigger.ID,
|
||||||
|
}
|
||||||
|
|
||||||
|
defaultResults = append(defaultResults, shuffle.ActionResult{
|
||||||
|
Action: curaction,
|
||||||
|
ExecutionId: workflowExecution.ExecutionId,
|
||||||
|
Authorization: workflowExecution.Authorization,
|
||||||
|
Result: "Skipped because it's not under the startnode",
|
||||||
|
StartedAt: 0,
|
||||||
|
CompletedAt: 0,
|
||||||
|
Status: "SKIPPED",
|
||||||
|
})
|
||||||
|
} else {
|
||||||
|
// Replaces trigger with the subflow
|
||||||
|
//if trigger.AppName == "Shuffle Workflow" {
|
||||||
|
// replaceActions := false
|
||||||
|
// workflowAction := ""
|
||||||
|
// for _, param := range trigger.Parameters {
|
||||||
|
// if param.Name == "argument" && !strings.Contains(param.Value, ".#") {
|
||||||
|
// replaceActions = true
|
||||||
|
// }
|
||||||
|
|
||||||
|
// if param.Name == "startnode" {
|
||||||
|
// workflowAction = param.Value
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
|
||||||
|
// if replaceActions {
|
||||||
|
// replacementNodes, newBranches, lastnode := shuffle.GetReplacementNodes(ctx, workflowExecution, trigger, trigger.Label)
|
||||||
|
// log.Printf("REPLACEMENTS: %d, %d", len(replacementNodes), len(newBranches))
|
||||||
|
// if len(replacementNodes) > 0 {
|
||||||
|
// for _, action := range replacementNodes {
|
||||||
|
// found := false
|
||||||
|
|
||||||
|
// for subActionIndex, subaction := range newActions {
|
||||||
|
// if subaction.ID == action.ID {
|
||||||
|
// found = true
|
||||||
|
// //newActions[subActionIndex].Name = action.Name
|
||||||
|
// newActions[subActionIndex].Label = action.Label
|
||||||
|
// break
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
|
||||||
|
// if !found {
|
||||||
|
// action.SubAction = true
|
||||||
|
// newActions = append(newActions, action)
|
||||||
|
// }
|
||||||
|
|
||||||
|
// // Check if it's already set to have a value
|
||||||
|
// for resultIndex, result := range defaultResults {
|
||||||
|
// if result.Action.ID == action.ID {
|
||||||
|
// defaultResults = append(defaultResults[:resultIndex], defaultResults[resultIndex+1:]...)
|
||||||
|
// break
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
|
||||||
|
// for _, branch := range newBranches {
|
||||||
|
// workflowExecution.Workflow.Branches = append(workflowExecution.Workflow.Branches, branch)
|
||||||
|
// }
|
||||||
|
|
||||||
|
// // Append branches:
|
||||||
|
// // parent -> new inner node (FIRST one)
|
||||||
|
// for branchIndex, branch := range workflowExecution.Workflow.Branches {
|
||||||
|
// if branch.DestinationID == trigger.ID {
|
||||||
|
// log.Printf("REPLACE DESTINATION WITH %s!!", workflowAction)
|
||||||
|
// workflowExecution.Workflow.Branches[branchIndex].DestinationID = workflowAction
|
||||||
|
// }
|
||||||
|
|
||||||
|
// if branch.SourceID == trigger.ID {
|
||||||
|
// log.Printf("REPLACE SOURCE WITH LASTNODE %s!!", lastnode)
|
||||||
|
// workflowExecution.Workflow.Branches[branchIndex].SourceID = lastnode
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
|
||||||
|
// // Remove the trigger
|
||||||
|
// removeTriggers = append(removeTriggers, workflowExecution.Workflow.Triggers[triggerIndex].ID)
|
||||||
|
// }
|
||||||
|
|
||||||
|
// log.Printf("NEW ACTION LENGTH %d, RESULT: %d, Triggers: %d, BRANCHES: %d", len(newActions), len(defaultResults), len(workflowExecution.Workflow.Triggers), len(workflowExecution.Workflow.Branches))
|
||||||
|
// }
|
||||||
|
//}
|
||||||
|
_ = triggerIndex
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
//newTriggers := []shuffle.Trigger{}
|
||||||
|
//for _, trigger := range workflowExecution.Workflow.Triggers {
|
||||||
|
// found := false
|
||||||
|
// for _, triggerId := range removeTriggers {
|
||||||
|
// if trigger.ID == triggerId {
|
||||||
|
// found = true
|
||||||
|
// break
|
||||||
|
// }
|
||||||
|
// }
|
||||||
|
|
||||||
|
// if found {
|
||||||
|
// log.Printf("[WARNING] Removed trigger %s during execution", trigger.ID)
|
||||||
|
// continue
|
||||||
|
// }
|
||||||
|
|
||||||
|
// newTriggers = append(newTriggers, trigger)
|
||||||
|
//}
|
||||||
|
//workflowExecution.Workflow.Triggers = newTriggers
|
||||||
|
_ = removeTriggers
|
||||||
|
|
||||||
|
if !startFound {
|
||||||
|
if len(workflowExecution.Start) == 0 && len(workflowExecution.Workflow.Start) > 0 {
|
||||||
|
workflowExecution.Start = workflow.Start
|
||||||
|
} else if len(workflowExecution.Workflow.Actions) > 0 {
|
||||||
|
workflowExecution.Start = workflowExecution.Workflow.Actions[0].ID
|
||||||
|
} else {
|
||||||
|
log.Printf("[ERROR] Startnode %s doesn't exist!!", workflowExecution.Start)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Workflow action %s doesn't exist in workflow", workflowExecution.Start), errors.New(fmt.Sprintf(`Workflow start node "%s" doesn't exist. Exiting!`, workflowExecution.Start))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
//log.Printf("EXECUTION START: %s", workflowExecution.Start)
|
||||||
|
|
||||||
|
// Verification for execution environments
|
||||||
|
workflowExecution.Results = defaultResults
|
||||||
|
workflowExecution.Workflow.Actions = newActions
|
||||||
|
onpremExecution := true
|
||||||
|
environments := []string{}
|
||||||
|
|
||||||
|
if len(workflowExecution.ExecutionOrg) == 0 && len(workflow.ExecutingOrg.Id) > 0 {
|
||||||
|
workflowExecution.ExecutionOrg = workflow.ExecutingOrg.Id
|
||||||
|
}
|
||||||
|
|
||||||
|
var allEnvs []shuffle.Environment
|
||||||
|
if len(workflowExecution.ExecutionOrg) > 0 {
|
||||||
|
//log.Printf("[INFO] Executing ORG: %s", workflowExecution.ExecutionOrg)
|
||||||
|
|
||||||
|
allEnvironments, err := shuffle.GetEnvironments(ctx, workflowExecution.ExecutionOrg)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("Failed finding environments: %s", err)
|
||||||
|
return shuffle.WorkflowExecution{}, fmt.Sprintf("Workflow environments not found for this org"), errors.New(fmt.Sprintf("Workflow environments not found for this org"))
|
||||||
|
}
|
||||||
|
|
||||||
|
for _, curenv := range allEnvironments {
|
||||||
|
if curenv.Archived {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
|
||||||
|
allEnvs = append(allEnvs, curenv)
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
log.Printf("[ERROR] No org identified for execution of %s. Returning", workflowExecution.Workflow.ID)
|
||||||
|
return shuffle.WorkflowExecution{}, "No org identified for execution", errors.New("No org identified for execution")
|
||||||
|
}
|
||||||
|
|
||||||
|
if len(allEnvs) == 0 {
|
||||||
|
log.Printf("[ERROR] No active environments found for org: %s", workflowExecution.ExecutionOrg)
|
||||||
|
return shuffle.WorkflowExecution{}, "No active environments found", errors.New(fmt.Sprintf("No active env found for org %s", workflowExecution.ExecutionOrg))
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if the actions are children of the startnode?
|
||||||
|
imageNames := []string{}
|
||||||
|
cloudExec := false
|
||||||
|
for _, action := range workflowExecution.Workflow.Actions {
|
||||||
|
// Verify if the action environment exists and append
|
||||||
|
found := false
|
||||||
|
for _, env := range allEnvs {
|
||||||
|
if env.Name == action.Environment {
|
||||||
|
found = true
|
||||||
|
|
||||||
|
if env.Type == "cloud" {
|
||||||
|
cloudExec = true
|
||||||
|
} else if env.Type == "onprem" {
|
||||||
|
onpremExecution = true
|
||||||
|
} else {
|
||||||
|
log.Printf("[ERROR] No handler for environment type %s", env.Type)
|
||||||
|
return shuffle.WorkflowExecution{}, "No active environments found", errors.New(fmt.Sprintf("No handler for environment type %s", env.Type))
|
||||||
|
}
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if !found {
|
||||||
|
log.Printf("[ERROR] Couldn't find environment %s. Maybe it's inactive?", action.Environment)
|
||||||
|
return shuffle.WorkflowExecution{}, "Couldn't find the environment", errors.New(fmt.Sprintf("Couldn't find env %s in org %s", action.Environment, workflowExecution.ExecutionOrg))
|
||||||
|
}
|
||||||
|
|
||||||
|
found = false
|
||||||
|
for _, env := range environments {
|
||||||
|
if env == action.Environment {
|
||||||
|
|
||||||
|
found = true
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if the app exists?
|
||||||
|
newName := action.AppName
|
||||||
|
newName = strings.ReplaceAll(newName, " ", "-")
|
||||||
|
imageNames = append(imageNames, fmt.Sprintf("%s:%s_%s", baseDockerName, newName, action.AppVersion))
|
||||||
|
|
||||||
|
if !found {
|
||||||
|
environments = append(environments, action.Environment)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
err = imageCheckBuilder(imageNames)
|
||||||
|
if err != nil {
|
||||||
|
log.Printf("[ERROR] Failed building the required images from %#v: %s", imageNames, err)
|
||||||
return shuffle.WorkflowExecution{}, "Failed building missing Docker images", err
|
return shuffle.WorkflowExecution{}, "Failed building missing Docker images", err
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2867,7 +3539,7 @@ func executeSingleAction(resp http.ResponseWriter, request *http.Request) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
resp.WriteHeader(200)
|
resp.WriteHeader(200)
|
||||||
resp.Write(returnBytes)
|
resp.Write([]byte(returnBytes))
|
||||||
}
|
}
|
||||||
|
|
||||||
// Onlyname is used to
|
// Onlyname is used to
|
||||||
|
|||||||
@@ -1,10 +1,9 @@
|
|||||||
{
|
{
|
||||||
"name": "shuffler",
|
"name": "shuffler",
|
||||||
"homepage": "https://shuffler.io",
|
"homepage": "https://shuffler.io",
|
||||||
"version": "1.1.4",
|
"version": "1.2.0",
|
||||||
"private": true,
|
"private": true,
|
||||||
"dependencies": {
|
"dependencies": {
|
||||||
"@babel/core": "^7.15.8",
|
|
||||||
"@codemirror/commands": "^6.2.2",
|
"@codemirror/commands": "^6.2.2",
|
||||||
"@emotion/is-prop-valid": "^1.1.1",
|
"@emotion/is-prop-valid": "^1.1.1",
|
||||||
"@emotion/react": "^11.7.0",
|
"@emotion/react": "^11.7.0",
|
||||||
@@ -23,7 +22,6 @@
|
|||||||
"@uiw/react-codemirror": "^3.2.1",
|
"@uiw/react-codemirror": "^3.2.1",
|
||||||
"@use-it/interval": "^1.0.0",
|
"@use-it/interval": "^1.0.0",
|
||||||
"algoliasearch": "^4.13.1",
|
"algoliasearch": "^4.13.1",
|
||||||
"babel-eslint": "^10.1.0",
|
|
||||||
"class-transformer": "^0.4.0",
|
"class-transformer": "^0.4.0",
|
||||||
"create-react-app": "^4.0.3",
|
"create-react-app": "^4.0.3",
|
||||||
"cytoscape": "^3.15.1",
|
"cytoscape": "^3.15.1",
|
||||||
@@ -79,14 +77,12 @@
|
|||||||
"shellwords": "^0.1.1",
|
"shellwords": "^0.1.1",
|
||||||
"simplebar": "^4.2.3",
|
"simplebar": "^4.2.3",
|
||||||
"styled-components": "^4.4.0",
|
"styled-components": "^4.4.0",
|
||||||
"webpack": "^4.44.2",
|
|
||||||
"websocket": "^1.0.30",
|
|
||||||
"yaml": "^1.7.2",
|
"yaml": "^1.7.2",
|
||||||
"yamljs": "^0.3.0",
|
"yamljs": "^0.3.0",
|
||||||
"zone.js": "~0.11.4"
|
"zone.js": "~0.11.4"
|
||||||
},
|
},
|
||||||
"scripts": {
|
"scripts": {
|
||||||
"start": "PORT=3000 react-scripts start",
|
"start": "HTTPS=false&&PORT=3000 react-scripts --openssl-legacy-provider start",
|
||||||
"build": "react-scripts build",
|
"build": "react-scripts build",
|
||||||
"test": "react-scripts test",
|
"test": "react-scripts test",
|
||||||
"eject": "react-scripts eject",
|
"eject": "react-scripts eject",
|
||||||
@@ -110,5 +106,9 @@
|
|||||||
"devDependencies": {
|
"devDependencies": {
|
||||||
"prettier": "2.4.1",
|
"prettier": "2.4.1",
|
||||||
"promise-window": "^1.2.1"
|
"promise-window": "^1.2.1"
|
||||||
|
"@babel/core": "^7.15.8",
|
||||||
|
"babel-eslint": "^10.1.0",
|
||||||
|
"webpack": "^4.44.2",
|
||||||
|
"@babel/plugin-proposal-private-property-in-object": "^7.21.11"
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
|
Before Width: | Height: | Size: 1.1 KiB After Width: | Height: | Size: 1.1 KiB |
|
Before Width: | Height: | Size: 217 KiB After Width: | Height: | Size: 217 KiB |
|
Before Width: | Height: | Size: 14 KiB After Width: | Height: | Size: 14 KiB |
|
Before Width: | Height: | Size: 20 KiB After Width: | Height: | Size: 20 KiB |
|
Before Width: | Height: | Size: 449 KiB After Width: | Height: | Size: 449 KiB |
|
Before Width: | Height: | Size: 431 KiB After Width: | Height: | Size: 431 KiB |
|
Before Width: | Height: | Size: 6.2 KiB After Width: | Height: | Size: 6.2 KiB |
|
Before Width: | Height: | Size: 9.6 KiB After Width: | Height: | Size: 9.6 KiB |
|
Before Width: | Height: | Size: 58 KiB After Width: | Height: | Size: 58 KiB |
|
Before Width: | Height: | Size: 2.0 KiB After Width: | Height: | Size: 2.0 KiB |
|
Before Width: | Height: | Size: 5.7 KiB After Width: | Height: | Size: 5.7 KiB |