Added initial opensearch setup

This commit is contained in:
frikky
2021-05-21 20:43:29 +02:00
parent 61416addaf
commit 4756d417b3
13 changed files with 180 additions and 584 deletions
+1 -1
View File
@@ -1,6 +1,6 @@
#!/bin/bash
NAME=shuffle-app_sdk
VERSION=0.8.82
VERSION=0.8.90
docker rmi docker.pkg.github.com/frikky/shuffle/$NAME:$VERSION --force
docker build . -t frikky/shuffle:app_sdk -t frikky/$NAME:$VERSION -t docker.pkg.github.com/frikky/shuffle/$NAME:$VERSION -t ghcr.io/frikky/$NAME:$VERSION
+1 -1
View File
@@ -24,7 +24,7 @@ require (
github.com/elastic/go-elasticsearch/v7 v7.12.0 // indirect
github.com/elastic/go-elasticsearch/v8 v8.0.0-20210519083322-55daf7425ecb // indirect
github.com/frikky/kin-openapi v0.39.0
github.com/frikky/shuffle-shared v0.0.46
github.com/frikky/shuffle-shared v0.0.47
github.com/fsouza/go-dockerclient v1.7.2 // indirect
github.com/ghodss/yaml v1.0.0
github.com/go-git/go-billy/v5 v5.0.0
+2
View File
@@ -151,6 +151,8 @@ github.com/frikky/shuffle-shared v0.0.38 h1:OZSwU1HDOaPzdlG1s77svgXJKzlNewM6GjeH
github.com/frikky/shuffle-shared v0.0.38/go.mod h1:H7SqOta/EAYnfYuWzwzYSh/oWfF0kgnuaJTQNKQBvoQ=
github.com/frikky/shuffle-shared v0.0.40 h1:H0au2np5xSy9mZEUWN+a29IORk+YP95FipENBD7iJRw=
github.com/frikky/shuffle-shared v0.0.40/go.mod h1:H7SqOta/EAYnfYuWzwzYSh/oWfF0kgnuaJTQNKQBvoQ=
github.com/frikky/shuffle-shared v0.0.46 h1:L52pyEVKZujM136qzD7LnarJIq0iXyfpDqhGrJkE7ik=
github.com/frikky/shuffle-shared v0.0.46/go.mod h1:BknTfpun3qte5bumR3OqQHf9XWPIsyj8woiXCjIlbBc=
github.com/fsouza/go-dockerclient v1.7.2 h1:bBEAcqLTkpq205jooP5RVroUKiVEWgGecHyeZc4OFjo=
github.com/fsouza/go-dockerclient v1.7.2/go.mod h1:+ugtMCVRwnPfY7d8/baCzZ3uwB0BrG5DB8OzbtxaRz8=
github.com/getkin/kin-openapi v0.8.0 h1:a6TQjTqwkyscC4/hShJX7WhCVE+4bi9lzw61XHQW5hE=
+135 -36
View File
@@ -780,6 +780,37 @@ func handleRegister(resp http.ResponseWriter, request *http.Request) {
}
} else {
log.Printf("[WARNING] Couldn't find an org to attach to. Create?")
orgSetupName := "default"
orgId := uuid.NewV4().String()
newOrg := shuffle.Org{
Name: orgSetupName,
Id: orgId,
Org: orgSetupName,
Users: []shuffle.User{},
Roles: []string{"admin", "user"},
CloudSync: false,
}
err = shuffle.SetOrg(ctx, newOrg, orgId)
if err != nil {
log.Printf("Failed setting organization: %s", err)
} else {
log.Printf("Successfully created the default org!")
item := shuffle.Environment{
Name: "Shuffle",
Type: "onprem",
OrgId: orgId,
Default: true,
Id: uuid.NewV4().String(),
}
err = shuffle.SetEnvironment(ctx, &item)
if err != nil {
log.Printf("[WARNING] Failed setting up new environment for new org")
}
}
}
}
@@ -983,6 +1014,64 @@ type passwordReset struct {
Reference string `json:"reference"`
}
// This might be... a bit off, but that's fine :)
// This might also be stupid, as we want timelines and such
// Anyway, these are super basic stupid stats.
func increaseStatisticsField(ctx context.Context, fieldname, id string, amount int64, orgId string) error {
// 1. Get current stats
// 2. Increase field(s)
// 3. Put new stats
statisticsId := "global_statistics"
nameKey := fieldname
key := datastore.NameKey(statisticsId, nameKey, nil)
statisticsItem := shuffle.StatisticsItem{}
newData := shuffle.StatisticsData{
Timestamp: int64(time.Now().Unix()),
Amount: amount,
Id: id,
}
if err := dbclient.Get(ctx, key, &statisticsItem); err != nil {
// Should init
if strings.Contains(fmt.Sprintf("%s", err), "entity") {
statisticsItem = shuffle.StatisticsItem{
Total: amount,
OrgId: orgId,
Fieldname: fieldname,
Data: []shuffle.StatisticsData{
newData,
},
}
if _, err := dbclient.Put(ctx, key, &statisticsItem); err != nil {
log.Printf("Error setting base stats: %s", err)
return err
}
return nil
}
//log.Printf("STATSERR: %s", err)
return err
}
statisticsItem.Total += amount
statisticsItem.Data = append(statisticsItem.Data, newData)
// New struct, to not add body, author etc
// FIXME - reintroduce
//if _, err := dbclient.Put(ctx, key, &statisticsItem); err != nil {
// log.Printf("Error stats to %s: %s", fieldname, err)
// return err
//}
//log.Printf("Stats: %#v", statisticsItem)
return nil
}
// FIXME - forward this to emails or whatever CRM system in use
func handleContact(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
@@ -3458,7 +3547,7 @@ func verifySwagger(resp http.ResponseWriter, request *http.Request) {
// Hint: Save API.id somewhere, and use newmd5 to save latest version
err = shuffle.SetOpenApiDatastore(ctx, newmd5, parsed)
if err != nil {
log.Printf("[ERROR] Failed saving to datastore: %s", err)
log.Printf("[ERROR] Failed saving app %s to database: %s", newmd5, err)
resp.WriteHeader(500)
resp.Write([]byte(fmt.Sprintf(`{"success": true, "reason": "%"}`, err)))
}
@@ -3927,7 +4016,7 @@ func runInitEs(ctx context.Context) {
log.Printf("ORGS: %d", len(activeOrgs))
if err != nil {
if fmt.Sprintf("%s", err) == "EOF" {
time.Sleep(5 * time.Second)
time.Sleep(7 * time.Second)
runInitEs(ctx)
return
}
@@ -3941,37 +4030,37 @@ func runInitEs(ctx context.Context) {
if len(activeOrgs) == 0 {
log.Printf(`No orgs. Setting NEW org "default"`)
orgSetupName := "default"
orgId := uuid.NewV4().String()
newOrg := shuffle.Org{
Name: orgSetupName,
Id: orgId,
Org: orgSetupName,
Users: []shuffle.User{},
Roles: []string{"admin", "user"},
CloudSync: false,
}
//orgSetupName := "default"
//orgId := uuid.NewV4().String()
//newOrg := shuffle.Org{
// Name: orgSetupName,
// Id: orgId,
// Org: orgSetupName,
// Users: []shuffle.User{},
// Roles: []string{"admin", "user"},
// CloudSync: false,
//}
err = shuffle.SetOrg(ctx, newOrg, orgId)
if err != nil {
log.Printf("Failed setting organization: %s", err)
} else {
log.Printf("Successfully created the default org!")
setUsers = true
}
//err = shuffle.SetOrg(ctx, newOrg, orgId)
//if err != nil {
// log.Printf("Failed setting organization: %s", err)
//} else {
// log.Printf("Successfully created the default org!")
// setUsers = true
item := shuffle.Environment{
Name: "Shuffle",
Type: "onprem",
OrgId: orgId,
Default: true,
Id: uuid.NewV4().String(),
}
// item := shuffle.Environment{
// Name: "Shuffle",
// Type: "onprem",
// OrgId: orgId,
// Default: true,
// Id: uuid.NewV4().String(),
// }
err = shuffle.SetEnvironment(ctx, &item)
if err != nil {
log.Printf("[WARNING] Failed setting up new environment for new org")
}
// err = shuffle.SetEnvironment(ctx, &item)
// if err != nil {
// log.Printf("[WARNING] Failed setting up new environment for new org")
// }
//}
} else {
log.Printf("[DEBUG] There are %d org(s).", len(activeOrgs))
@@ -4120,7 +4209,7 @@ func runInitEs(ctx context.Context) {
if err != nil && len(workflowapps) == 0 {
log.Printf("[WARNING] Failed getting apps (runInit): %s", err)
} else if err == nil {
log.Printf("Downloading default workflow apps")
log.Printf("[DEBUG] Downloading default apps")
fs := memfs.New()
storer := memory.NewStorage()
@@ -5237,9 +5326,11 @@ func handleCloudSetup(resp http.ResponseWriter, request *http.Request) {
// 1. Find environment
// 2. If cloud env found, enable it (un-archive)
// 3. If it doesn't create it
var environments []shuffle.Environment
q := datastore.NewQuery("Environments").Filter("org_id =", org.Id)
_, err = dbclient.GetAll(ctx, q, &environments)
//var environments []shuffle.Environment
//q := datastore.NewQuery("Environments").Filter("org_id =", org.Id)
//_, err = dbclient.GetAll(ctx, q, &environments)
environments, err := shuffle.GetEnvironments(ctx, org.Id)
if err == nil {
// Don't disable, this will be deleted entirely
@@ -5573,11 +5664,19 @@ func initHandlers() {
panic(fmt.Sprintf("[DEBUG] Database client for ELASTICSEARCH error during init: %s", err))
}
elasticConfig := "elasticsearch"
if strings.ToLower(os.Getenv("SHUFFLE_ELASTIC")) == "false" {
elasticConfig = ""
}
_ = shuffle.RunInit(*dbclient, *es, storage.Client{}, gceProject, "onprem", true, "elasticsearch")
log.Printf("[DEBUG] Finished Shuffle database init")
//go runInit(ctx)
go runInitEs(ctx)
if elasticConfig == "elasticsearch" {
go runInitEs(ctx)
} else {
go runInit(ctx)
}
r := mux.NewRouter()
r.HandleFunc("/api/v1/_ah/health", healthCheckHandler)
+2 -2
View File
@@ -802,7 +802,7 @@ func handleDeleteOutlookSub(resp http.ResponseWriter, request *http.Request) {
ctx := context.Background()
workflow, err := shuffle.GetWorkflow(ctx, workflowId)
if err != nil {
log.Printf("Failed getting the workflow locally (delete outlook): %s", err)
log.Printf("[WARNING] Failed getting the workflow locally (delete outlook): %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
@@ -863,7 +863,7 @@ func createOutlookSub(resp http.ResponseWriter, request *http.Request) {
ctx := context.Background()
workflow, err := shuffle.GetWorkflow(ctx, workflowId)
if err != nil {
log.Printf("Failed getting the workflow locally (outlook sub): %s", err)
log.Printf("[WARNING] Failed getting the workflow locally (outlook sub): %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
+4 -512
View File
@@ -22,7 +22,6 @@ import (
"github.com/docker/docker/api/types"
"github.com/docker/docker/client"
"cloud.google.com/go/datastore"
scheduler "cloud.google.com/go/scheduler/apiv1"
gyaml "github.com/ghodss/yaml"
"github.com/h2non/filetype"
@@ -43,8 +42,6 @@ import (
//"google.golang.org/appengine/memcache"
//"cloud.google.com/go/firestore"
// "google.golang.org/api/option"
"google.golang.org/api/iterator"
)
var localBase = "http://localhost:5001"
@@ -474,65 +471,6 @@ var scheduledOrgs = map[string]*newscheduler.Job{}
// SuccessExamples []string `json:"success_examples" datastore:"success_examples,noindex"`
// FailureExamples []string `json:"failure_examples" datastore:"failure_examples,noindex"`
//}
// This might be... a bit off, but that's fine :)
// This might also be stupid, as we want timelines and such
// Anyway, these are super basic stupid stats.
func increaseStatisticsField(ctx context.Context, fieldname, id string, amount int64, orgId string) error {
// 1. Get current stats
// 2. Increase field(s)
// 3. Put new stats
statisticsId := "global_statistics"
nameKey := fieldname
key := datastore.NameKey(statisticsId, nameKey, nil)
statisticsItem := shuffle.StatisticsItem{}
newData := shuffle.StatisticsData{
Timestamp: int64(time.Now().Unix()),
Amount: amount,
Id: id,
}
if err := dbclient.Get(ctx, key, &statisticsItem); err != nil {
// Should init
if strings.Contains(fmt.Sprintf("%s", err), "entity") {
statisticsItem = shuffle.StatisticsItem{
Total: amount,
OrgId: orgId,
Fieldname: fieldname,
Data: []shuffle.StatisticsData{
newData,
},
}
if _, err := dbclient.Put(ctx, key, &statisticsItem); err != nil {
log.Printf("Error setting base stats: %s", err)
return err
}
return nil
}
//log.Printf("STATSERR: %s", err)
return err
}
statisticsItem.Total += amount
statisticsItem.Data = append(statisticsItem.Data, newData)
// New struct, to not add body, author etc
// FIXME - reintroduce
//if _, err := dbclient.Put(ctx, key, &statisticsItem); err != nil {
// log.Printf("Error stats to %s: %s", fieldname, err)
// return err
//}
//log.Printf("Stats: %#v", statisticsItem)
return nil
}
// Frequency = cronjob OR minutes between execution
func createSchedule(ctx context.Context, scheduleId, workflowId, name, startNode, frequency, orgId string, body []byte) error {
var err error
@@ -1169,98 +1107,6 @@ func handleExecutionStatistics(execution shuffle.WorkflowExecution) {
}
}
func getWorkflows(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
return
}
user, err := shuffle.HandleApiAuthentication(resp, request)
if err != nil {
log.Printf("Api authentication failed in getworkflows: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
//memcacheName := fmt.Sprintf("%s_workflows", user.Username)
ctx := context.Background()
//if item, err := memcache.Get(ctx, memcacheName); err == memcache.ErrCacheMiss {
// // Not in cache
// //log.Printf("Workflows not in cache.")
//} else if err != nil {
// log.Printf("Error getting item: %v", err)
//} else {
// // FIXME - verify if value is ok? Can unmarshal etc.
// resp.WriteHeader(200)
// resp.Write(item.Value)
// return
//}
// With user, do a search for workflows with user or user's org attached
q := datastore.NewQuery("workflow").Filter("owner =", user.Id)
if user.Role == "admin" {
q = datastore.NewQuery("workflow").Filter("org_id =", user.ActiveOrg.Id)
log.Printf("[INFO] Getting workflows (ADMIN) for organization %s", user.ActiveOrg.Id)
}
q = q.Order("-edited")
var workflows []shuffle.Workflow
_, err = dbclient.GetAll(ctx, q, &workflows)
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "ResourceExhausted") {
q = q.Limit(35)
_, err = dbclient.GetAll(ctx, q, &workflows)
if err != nil {
log.Printf("Failed getting workflows for user %s: %s (0)", user.Username, err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
} else {
log.Printf("Failed getting workflows for user %s: %s (1)", user.Username, err)
//shuffle.DeleteKey(ctx, "workflow", "5694357e-8063-4580-8529-301cc72df951")
//log.Printf("Workflows: %#v", workflows)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
}
if len(workflows) == 0 {
resp.WriteHeader(200)
resp.Write([]byte("[]"))
return
}
newjson, err := json.Marshal(workflows)
if err != nil {
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed unpacking workflows"}`)))
return
}
//item := &memcache.Item{
// Key: memcacheName,
// Value: newjson,
// Expiration: time.Minute * 10,
//}
//if err := memcache.Add(ctx, item); err == memcache.ErrNotStored {
// if err := memcache.Set(ctx, item); err != nil {
// log.Printf("Error setting item: %v", err)
// }
//} else if err != nil {
// log.Printf("Error adding item: %v", err)
//} else {
// //log.Printf("Set cache for %s", item.Key)
//}
resp.WriteHeader(200)
resp.Write(newjson)
}
func deleteWorkflow(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
@@ -1297,7 +1143,7 @@ func deleteWorkflow(resp http.ResponseWriter, request *http.Request) {
ctx := context.Background()
workflow, err := shuffle.GetWorkflow(ctx, fileId)
if err != nil {
log.Printf("Failed getting the workflow locally (delete workflow): %s", err)
log.Printf("[WARNING] Failed getting workflow %s locally (delete workflow): %s", fileId, err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
@@ -1340,7 +1186,6 @@ func deleteWorkflow(resp http.ResponseWriter, request *http.Request) {
}
// FIXME - maybe delete workflow executions
log.Printf("[INFO] Should have deleted workflow %s", fileId)
err = shuffle.DeleteKey(ctx, "workflow", fileId)
if err != nil {
log.Printf("Failed deleting key %s", fileId)
@@ -1348,15 +1193,14 @@ func deleteWorkflow(resp http.ResponseWriter, request *http.Request) {
resp.Write([]byte(`{"success": false, "reason": "Failed deleting key"}`))
return
}
log.Printf("[INFO] Should have deleted workflow %s", fileId)
//err = increaseStatisticsField(ctx, "total_workflows", fileId, -1, workflow.OrgId)
//if err != nil {
// log.Printf("Failed to increase total workflows: %s", err)
//}
//memcacheName := fmt.Sprintf("%s_%s", user.Username, fileId)
//memcache.Delete(ctx, memcacheName)
//memcacheName = fmt.Sprintf("%s_workflows", user.Username)
//memcache.Delete(ctx, memcacheName)
cacheKey := fmt.Sprintf("%s_workflows", user.Id)
shuffle.DeleteCache(ctx, cacheKey)
resp.WriteHeader(200)
resp.Write([]byte(`{"success": true}`))
@@ -1401,45 +1245,6 @@ func getWorkflowLocal(fileId string, request *http.Request) ([]byte, error) {
//// New execution with firestore
func cleanupExecutions(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
return
}
user, err := shuffle.HandleApiAuthentication(resp, request)
if err != nil {
log.Printf("[INFO] Api authentication failed in cleanup executions: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false, "message": "Not authenticated"}`))
return
}
if user.Role != "admin" {
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false, "message": "Insufficient permissions"}`))
return
}
ctx := context.Background()
// Removes three months from today
timestamp := int64(time.Now().AddDate(0, -2, 0).Unix())
log.Println(timestamp)
q := datastore.NewQuery("workflowexecution").Filter("started_at <", timestamp)
var workflowExecutions []shuffle.WorkflowExecution
_, err = dbclient.GetAll(ctx, q, &workflowExecutions)
if err != nil {
log.Printf("Error getting workflowexec (cleanup): %s", err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting all workflowexecutions"}`)))
return
}
resp.WriteHeader(200)
resp.Write([]byte(`{"success": true}`))
}
func handleExecution(id string, workflow shuffle.Workflow, request *http.Request) (shuffle.WorkflowExecution, string, error) {
ctx := context.Background()
if workflow.ID == "" || workflow.ID != id {
@@ -2851,168 +2656,6 @@ func setExampleresult(ctx context.Context, result shuffle.AppExecutionExample) e
return nil
}
// FIXME: Not suitable for cloud right now :O
func deleteWorkflowApp(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
return
}
user, userErr := shuffle.HandleApiAuthentication(resp, request)
if userErr != nil {
log.Printf("Api authentication failed in edit workflow: %s", userErr)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
location := strings.Split(request.URL.String(), "/")
log.Printf("%#v", location)
var fileId string
if location[1] == "api" {
if len(location) <= 4 {
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
fileId = location[4]
}
ctx := context.Background()
log.Printf("ID: %s", fileId)
app, err := shuffle.GetApp(ctx, fileId, user)
if err != nil {
log.Printf("Error getting app (delete) %s: %s", fileId, err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
// FIXME - check whether it's in use and maybe restrict again for later?
// FIXME - actually delete other than private apps too..
private := false
if app.Downloaded && user.Role == "admin" {
log.Printf("[INFO] Deleting downloaded app (authenticated users can do this)")
} else if user.Id != app.Owner {
log.Printf("[WARNING] Wrong user (%s) for app %s (delete)", user.Username, app.Name)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
} else {
log.Printf("[WARNING] App to be deleted is private")
private = true
}
// FIXME: Make workflows track themself INSIDE apps, or with a reference
q := datastore.NewQuery("workflow").Filter("org_id = ", user.ActiveOrg.Id).Limit(30)
var workflows []shuffle.Workflow
_, err = dbclient.GetAll(ctx, q, &workflows)
if err != nil {
log.Printf("[WARNING] Failed getting related workflows for the app: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err)))
return
}
// Finds workflows using the app to set errors
// FIXME: this will be WAY too big for cloud :O
for _, workflow := range workflows {
found := false
newActions := []shuffle.Action{}
for _, action := range workflow.Actions {
if action.AppName == app.Name && action.AppVersion == app.AppVersion {
found = true
action.Errors = append(action.Errors, "App has been deleted")
action.IsValid = false
}
newActions = append(newActions, action)
}
if found {
workflow.IsValid = false
workflow.Errors = append(workflow.Errors, fmt.Sprintf("App %s_%s has been deleted", app.Name, app.AppVersion))
workflow.Actions = newActions
for _, trigger := range workflow.Triggers {
_ = trigger
//log.Printf("TRIGGER: %#v", trigger)
//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
//}
}
err = shuffle.SetWorkflow(ctx, workflow, workflow.ID)
if err != nil {
log.Printf("Failed setting workflow when deleting app: %s", err)
continue
} else {
log.Printf("Set %s (%s) to have errors", workflow.ID, workflow.Name)
}
}
}
//resp.WriteHeader(200)
//resp.Write([]byte(`{"success": true}`))
//return
// Not really deleting it, just removing from user cache
if private {
log.Printf("[INFO] Deleting private app")
var privateApps []shuffle.WorkflowApp
for _, item := range user.PrivateApps {
if item.ID == fileId {
continue
}
privateApps = append(privateApps, item)
}
user.PrivateApps = privateApps
err = shuffle.SetUser(ctx, &user, true)
if err != nil {
log.Printf("[ERROR] Failed removing %s app for user %s: %s", app.Name, user.Username, err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": true"}`)))
return
}
}
log.Printf("[INFO] Deleting public app")
err = shuffle.DeleteKey(ctx, "workflowapp", fileId)
if err != nil {
log.Printf("Failed deleting workflowapp")
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed deleting workflow app"}`)))
return
}
err = increaseStatisticsField(ctx, "total_apps_deleted", fileId, 1, user.ActiveOrg.Id)
if err != nil {
log.Printf("Failed to increase total apps loaded stats: %s", err)
}
cacheKey := fmt.Sprintf("workflowapps-sorted-100")
shuffle.DeleteCache(ctx, cacheKey)
cacheKey = fmt.Sprintf("workflowapps-sorted-500")
shuffle.DeleteCache(ctx, cacheKey)
//err = memcache.Delete(request.Context(), sessionToken)
resp.WriteHeader(200)
resp.Write([]byte(`{"success": true}`))
}
func getWorkflowApps(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
@@ -4433,157 +4076,6 @@ func setNewWorkflowApp(resp http.ResponseWriter, request *http.Request) {
resp.Write([]byte(fmt.Sprintf(`{"success": true}`)))
}
func getWorkflowExecutions(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)
if cors {
return
}
user, err := shuffle.HandleApiAuthentication(resp, request)
if err != nil {
log.Printf("Api authentication failed in getting specific workflow: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
location := strings.Split(request.URL.String(), "/")
var fileId string
if location[1] == "api" {
if len(location) <= 4 {
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
fileId = location[4]
}
if len(fileId) != 36 {
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false, "reason": "Workflow ID when getting workflow executions is not valid"}`))
return
}
ctx := context.Background()
workflow, err := shuffle.GetWorkflow(ctx, fileId)
if err != nil {
log.Printf("Failed getting the workflow %s locally (get executions): %s", fileId, err)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
// FIXME - have a check for org etc too..
if user.Id != workflow.Owner {
log.Printf("Wrong user (%s) for workflow %s (get execution)", user.Username, workflow.ID)
resp.WriteHeader(401)
resp.Write([]byte(`{"success": false}`))
return
}
// Query for the specifci workflowId
maxAmount := 30
q := datastore.NewQuery("workflowexecution").Filter("workflow_id =", fileId).Order("-started_at").Limit(maxAmount)
var workflowExecutions []shuffle.WorkflowExecution
_, err = dbclient.GetAll(ctx, q, &workflowExecutions)
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "ResourceExhausted") {
q = datastore.NewQuery("workflowexecution").Filter("workflow_id =", fileId).Order("-started_at").Limit(1)
/*
_, err = dbclient.GetAll(ctx, q, &workflowExecutions)
if err != nil {
log.Printf("Error getting workflowexec (2): %s", err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting all workflowexecutions for %s"}`, fileId)))
return
}
*/
cursorStr := ""
for {
it := dbclient.Run(ctx, q)
//_, err = it.Next(&app)
for {
var workflowExecution shuffle.WorkflowExecution
_, err := it.Next(&workflowExecution)
if err != nil {
break
}
workflowExecutions = append(workflowExecutions, workflowExecution)
}
//log.Printf("Len: %d", len(workflowExecutions))
if len(workflowExecutions) > maxAmount {
break
}
nextCursor, err := it.Cursor()
if err != iterator.Done && err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "ResourceExhausted") {
//log.Printf("NEXT!")
nextStr := fmt.Sprintf("%s", nextCursor)
if cursorStr == nextStr {
break
}
cursorStr = nextStr
continue
} else {
log.Printf("BREAK: %s", err)
break
}
}
if err != nil {
if strings.Contains(fmt.Sprintf("%s", err), "ResourceExhausted") {
log.Printf("[WARNING] Cursorerror in app grab WARNING: %s", err)
} else {
log.Printf("[ERROR] Cursorerror in app grab: %s", err)
break
}
} else {
//log.Printf("NEXTCURSOR: %s", nextCursor)
nextStr := fmt.Sprintf("%s", nextCursor)
if cursorStr == nextStr {
break
}
cursorStr = nextStr
q = q.Start(nextCursor)
//cursorStr = nextCursor
//break
}
}
} else {
log.Printf("Error getting workflowexec: %s", err)
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed getting all workflowexecutions for %s"}`, fileId)))
return
}
}
if len(workflowExecutions) == 0 {
resp.WriteHeader(200)
resp.Write([]byte("[]"))
return
}
newjson, err := json.Marshal(workflowExecutions)
if err != nil {
resp.WriteHeader(401)
resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "Failed unpacking workflow executions"}`)))
return
}
resp.WriteHeader(200)
resp.Write(newjson)
}
// Starts a new webhook
func handleStopHook(resp http.ResponseWriter, request *http.Request) {
cors := handleCors(resp, request)