diff --git a/backend/go-app/main.go b/backend/go-app/main.go index fbd54cc4..1c313be2 100644 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -42,6 +42,7 @@ import ( // Random xj "github.com/basgys/goxml2json" + newscheduler "github.com/carlescere/scheduler" gyaml "github.com/ghodss/yaml" "github.com/satori/go.uuid" "golang.org/x/crypto/bcrypt" @@ -231,10 +232,12 @@ type AppInfo struct { DestinationApp ScheduleApp `json:"destinationapp,omitempty" datastore:"destinationapp,noindex"` } -// Used for the api integrator -//Username string `datastore:"Username,noindex"` +// May 2020: Reused for onprem schedules - Id, Seconds, WorkflowId and argument type ScheduleOld struct { Id string `json:"id" datastore:"id"` + Seconds int `json:"seconds" datastore:"seconds"` + WorkflowId string `json:"workflow_id datastore:"workflow_id", ` + Argument string `json:"argument" datastore:"argument"` AppInfo AppInfo `json:"appinfo" datastore:"appinfo,noindex"` Finished bool `json:"finished" finished:"id"` BaseAppLocation string `json:"base_app_location" datastore:"baseapplocation,noindex"` @@ -5624,11 +5627,13 @@ func init() { log.Fatalf("DBclient error during init: %s", err) } + // Setting stats for backend starts (failure count as well) err = increaseStatisticsField(ctx, "backend_executions", "", 1) if err != nil { log.Printf("Failed increasing local stats: %s", err) } + // Gets environments and inits if it doesn't exist count, err := getEnvironmentCount() if count == 0 && err == nil { item := Environment{ @@ -5666,6 +5671,35 @@ func init() { iterateAppGithubFolders(fs, dir, "", "testing") } + // Gets schedules and starts them + schedules, err := getAllSchedules(ctx) + if err != nil { + log.Printf("Failed getting schedules during service init: %s", err) + } else { + log.Printf("Setting up %d schedule(s)", len(schedules)) + for _, schedule := range schedules { + job := func() { + request := &http.Request{ + Method: "POST", + Body: ioutil.NopCloser(strings.NewReader(schedule.Argument)), + } + + _, _, err := handleExecution(schedule.WorkflowId, Workflow{}, request) + if err != nil { + log.Printf("Failed to execute: %s", err) + } + } + + jobret, err := newscheduler.Every(schedule.Seconds).Seconds().NotImmediately().Run(job) + if err != nil { + log.Printf("Failed to schedule workflow: %s", err) + // FIXME: what now? lol:w + } + + scheduledJobs[schedule.Id] = jobret + } + } + log.Printf("Finished INIT") r := mux.NewRouter() diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 7ee828ba..ff124dad 100644 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -27,9 +27,8 @@ import ( "github.com/go-git/go-billy/v5/memfs" "github.com/go-git/go-git/v5" - "github.com/go-git/go-git/v5/storage/memory" - newscheduler "github.com/carlescere/scheduler" + "github.com/go-git/go-git/v5/storage/memory" //"github.com/gorilla/websocket" //"google.golang.org/appengine" //"google.golang.org/appengine/memcache" @@ -43,6 +42,7 @@ var baseEnvironment = "onprem" var cloudname = "cloud" var defaultLocation = "europe-west2" +var scheduledJobs = map[string]*newscheduler.Job{} // To test out firestore before potential merge var shuffleTestProject = "shuffle-test-258209" @@ -384,29 +384,35 @@ func getWorkflowQueue(ctx context.Context, id string) (ExecutionRequestWrapper, // Frequency = cronjob OR minutes between execution func createSchedule(ctx context.Context, scheduleId, workflowId, name, frequency string, body []byte) error { + var err error testSplit := strings.Split(frequency, "*") cronJob := "" + newfrequency := 0 + if len(testSplit) > 5 { cronJob = frequency } else { - newfrequency, err := strconv.Atoi(frequency) + newfrequency, err = strconv.Atoi(frequency) if err != nil { + log.Printf("Failed to parse time: %s", err) return err } - _ = newfrequency - //if int(newfrequency) < 60 { // cronJob = fmt.Sprintf("*/%s * * * *") //} else if int(newfrequency) < - log.Println("FIXME: SHOULD DO Frequency (minutes) to CRON") } - if len(cronJob) == 0 { + // Reverse. Can't handle CRON, only numbers + if len(cronJob) > 0 { return errors.New("cronJob isn't formatted correctly") } - log.Printf("CRON: %s, body: %s", cronJob, string(body)) + if newfrequency < 1 { + return errors.New("Frequency has to be more than 0") + } + + //log.Printf("CRON: %s, body: %s", cronJob, string(body)) // FIXME: // This may run multiple places if multiple servers, @@ -418,49 +424,42 @@ func createSchedule(ctx context.Context, scheduleId, workflowId, name, frequency } _, _, err := handleExecution(workflowId, Workflow{}, request) - if err == nil { + if err != nil { log.Printf("Failed to execute: %s", err) } } - // FIXME - Create a real schedule based on cron: - // 1. Parse the cron in a function to match this schedule - // 2. Make main init check for schedules that aren't running - _, err := newscheduler.Every(5).Seconds().NotImmediately().Run(job) + log.Printf("Starting frequency: %d", newfrequency) + jobret, err := newscheduler.Every(newfrequency).Seconds().NotImmediately().Run(job) if err != nil { log.Printf("Failed to schedule workflow: %s", err) return err } - return errors.New("ERROR!!") + //scheduledJobs = append(scheduledJobs, jobret) + scheduledJobs[scheduleId] = jobret - //log.Printf("REQUEST: %#v", executionRequest) + // Doesn't need running/not running. If stopped, we just delete it. + timeNow := int64(time.Now().Unix()) + schedule := ScheduleOld{ + Id: scheduleId, + WorkflowId: workflowId, + Argument: string(body), + Seconds: newfrequency, + CreationTime: timeNow, + LastModificationtime: timeNow, + LastRuntime: timeNow, + } - //req := &schedulerpb.CreateJobRequest{ - // Parent: fmt.Sprintf("projects/%s/locations/europe-west2", gceProject), - // Job: &schedulerpb.Job{ - // Name: fmt.Sprintf("projects/%s/locations/europe-west2/jobs/schedule_%s", gceProject, scheduleId), - // Schedule: cronJob, - // Description: name, - // Target: &schedulerpb.Job_HttpTarget{ - // HttpTarget: &schedulerpb.HttpTarget{ - // Uri: fmt.Sprintf("https://shuffler.io/api/v1/workflows/%s/execute", workflowId), - // HttpMethod: 1, - // Headers: map[string]string{ - // "Authorization": "", - // }, - // Body: body, - // }, - // }, - // }, - // // TODO: Fill request struct fields. - //} - //resp, err := c.CreateJob(ctx, req) - //if err != nil { - // log.Printf("%s", err) - // return err - //} - //_ = resp + err = setSchedule(ctx, schedule) + if err != nil { + log.Printf("Failed to set schedule: %s", err) + return err + } + + // FIXME - Create a real schedule based on cron: + // 1. Parse the cron in a function to match this schedule + // 2. Make main init check for schedules that aren't running return nil } @@ -1692,11 +1691,12 @@ func handleExecution(id string, workflow Workflow, request *http.Request) (Workf return WorkflowExecution{}, "Failed getting body", err } + // This one doesn't really matter. var execution ExecutionRequest err = json.Unmarshal(body, &execution) if err != nil { - log.Printf("Failed execution POST unmarshaling: %s", err) - return WorkflowExecution{}, "", err + //log.Printf("Failed execution POST unmarshaling - still continue: %s", err) + //return WorkflowExecution{}, "", err } if execution.Start == "" && len(body) > 0 { @@ -1708,7 +1708,7 @@ func handleExecution(id string, workflow Workflow, request *http.Request) (Workf workflowExecution.ExecutionArgument = execution.ExecutionArgument } - log.Printf("Execution data: %#v", execution) + //log.Printf("Execution data: %#v", execution) if len(execution.Start) == 36 { log.Printf("SHOULD START ON NODE %s", execution.Start) workflow.Start = execution.Start @@ -1916,7 +1916,7 @@ func handleExecution(id string, workflow Workflow, request *http.Request) (Workf executionRequestWrapper.Data = append(executionRequestWrapper.Data, executionRequest) } - log.Printf("Execution request: %#v", executionRequest) + //log.Printf("Execution request: %#v", executionRequest) err = setWorkflowQueue(ctx, executionRequestWrapper, environment) if err != nil { @@ -2086,7 +2086,117 @@ func stopSchedule(resp http.ResponseWriter, request *http.Request) { return } +func stopScheduleGCP(resp http.ResponseWriter, request *http.Request) { + cors := handleCors(resp, request) + if cors { + return + } + + user, err := handleApiAuthentication(resp, request) + if err != nil { + log.Printf("Api authentication failed in schedule workflow: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + location := strings.Split(request.URL.String(), "/") + + var fileId string + var scheduleId string + if location[1] == "api" { + if len(location) <= 6 { + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + fileId = location[4] + scheduleId = location[6] + } + + if len(fileId) != 36 { + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false, "reason": "Workflow ID to stop schedule is not valid"}`)) + return + } + + if len(scheduleId) != 36 { + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false, "reason": "Schedule ID not valid"}`)) + return + } + + ctx := context.Background() + workflow, err := getWorkflow(ctx, fileId) + if err != nil { + log.Printf("Failed getting the workflow locally: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + // FIXME - have a check for org etc too.. + // FIXME - admin check like this? idk + if user.Id != workflow.Owner && user.Role != "admin" && user.Role != "scheduler" { + log.Printf("Wrong user (%s) for workflow %s (stop schedule)", user.Username, workflow.ID) + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + if len(workflow.Actions) == 0 { + workflow.Actions = []Action{} + } + if len(workflow.Branches) == 0 { + workflow.Branches = []Branch{} + } + if len(workflow.Triggers) == 0 { + workflow.Triggers = []Trigger{} + } + if len(workflow.Errors) == 0 { + workflow.Errors = []string{} + } + + err = deleteSchedule(ctx, scheduleId) + if err != nil { + if strings.Contains(err.Error(), "Job not found") { + resp.WriteHeader(200) + resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) + } else { + resp.WriteHeader(401) + resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed stopping schedule"}`))) + } + return + } + + resp.WriteHeader(200) + resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) + return +} + func deleteSchedule(ctx context.Context, id string) error { + log.Printf("Should stop schedule %s!", id) + //newscheduler "github.com/carlescere/scheduler" + log.Printf("Schedules: %#v", scheduledJobs) + if value, exists := scheduledJobs[id]; exists { + log.Printf("STOP THIS ONE: %s", value) + // Looks like this does the trick? Hurr + value.Lock() + err := DeleteKey(ctx, "schedules", id) + if err != nil { + log.Printf("Failed to delete schedule: %s", err) + return err + } + } else { + // FIXME - allow it to kind of stop anyway? + return errors.New("Can't find the schedule.") + } + + return nil +} + +func deleteScheduleGCP(ctx context.Context, id string) error { c, err := scheduler.NewCloudSchedulerClient(ctx) if err != nil { log.Printf("%s", err) @@ -2208,13 +2318,7 @@ func scheduleWorkflow(resp http.ResponseWriter, request *http.Request) { return } - type tmp struct { - ExecutionArgument string `json:"execution_argument"` - } - - var tmpArg tmp - tmpArg.ExecutionArgument = schedule.ExecutionArgument - scheduleArg, err := json.Marshal(tmpArg) + scheduleArg, err := json.Marshal(schedule.ExecutionArgument) if err != nil { log.Printf("Failed scheduleArg marshal: %s", err) resp.WriteHeader(http.StatusInternalServerError) @@ -2222,6 +2326,8 @@ func scheduleWorkflow(resp http.ResponseWriter, request *http.Request) { return } + log.Printf("Schedulearg: %s", string(scheduleArg)) + err = createSchedule( ctx, schedule.Id, @@ -3315,7 +3421,7 @@ func getWorkflowExecutions(resp http.ResponseWriter, request *http.Request) { } // Query for the specifci workflowId - q := datastore.NewQuery("workflowexecution").Filter("workflow_id =", fileId) + q := datastore.NewQuery("workflowexecution").Filter("workflow_id =", fileId).Limit(50) var workflowExecutions []WorkflowExecution _, err = dbclient.GetAll(ctx, q, &workflowExecutions) if err != nil { @@ -3342,6 +3448,18 @@ func getWorkflowExecutions(resp http.ResponseWriter, request *http.Request) { resp.Write(newjson) } +func getAllSchedules(ctx context.Context) ([]ScheduleOld, error) { + var schedules []ScheduleOld + q := datastore.NewQuery("schedules") + + _, err := dbclient.GetAll(ctx, q, &schedules) + if err != nil { + return []ScheduleOld{}, err + } + + return schedules, nil +} + func getAllWorkflowApps(ctx context.Context) ([]WorkflowApp, error) { var allworkflowapps []WorkflowApp q := datastore.NewQuery("workflowapp") diff --git a/frontend/src/AngularWorkflow.js b/frontend/src/AngularWorkflow.js index 922e4240..7af0d0a9 100644 --- a/frontend/src/AngularWorkflow.js +++ b/frontend/src/AngularWorkflow.js @@ -80,7 +80,6 @@ const splitter = "|~|" //const referenceUrl = "https://shuffler.io/functions/webhooks/" //const referenceUrl = window.location.origin+"/api/v1/hooks/" -console.log(window.location) const AngularWorkflow = (props) => { const { globalUrl, isLoggedIn, isLoaded } = props; const referenceUrl = globalUrl+"/api/v1/hooks/" @@ -1261,11 +1260,11 @@ const AngularWorkflow = (props) => { //const submitSchedule = (id, name, frequency, executionArg) => { const submitSchedule = (trigger, triggerindex) => { - const cronSplit = workflow.triggers[triggerindex].parameters[0].value.split("*") - if (cronSplit.length <= 5 || cronSplit.length > 6) { - alert.error("Error: Bad cron, example run every 1 minute: */1 * * * *") - return - } + //const cronSplit = workflow.triggers[triggerindex].parameters[0].value.split("*") + //if (cronSplit.length <= 5 || cronSplit.length > 6) { + // alert.error("Error: Bad cron, example run every 1 minute: */1 * * * *") + // return + //} if (trigger.name.length <= 0) { alert.error("Error: name can't be empty") @@ -3892,7 +3891,7 @@ const AngularWorkflow = (props) => { if (Object.getOwnPropertyNames(selectedTrigger).length > 0 && workflow.triggers[selectedTriggerIndex] !== undefined) { if (workflow.triggers[selectedTriggerIndex].parameters === undefined || workflow.triggers[selectedTriggerIndex].parameters === null || workflow.triggers[selectedTriggerIndex].parameters.length === 0) { workflow.triggers[selectedTriggerIndex].parameters = [] - workflow.triggers[selectedTriggerIndex].parameters[0] = {"name": "cron", "value": "*/15 * * * *"} + workflow.triggers[selectedTriggerIndex].parameters[0] = {"name": "cron", "value": "120"} workflow.triggers[selectedTriggerIndex].parameters[1] = {"name": "execution_argument", "value": '{"example": {"json": "is cool"}}'} setWorkflow(workflow) } @@ -3953,7 +3952,7 @@ const AngularWorkflow = (props) => {