From 1d559a403aba36623d48f2ec6899aafafaa50792 Mon Sep 17 00:00:00 2001 From: frikky Date: Thu, 20 May 2021 20:57:52 +0200 Subject: [PATCH] Did major database migration to Opensearch --- backend/go-app/go.mod | 3 + backend/go-app/go.sum | 6 + backend/go-app/main.go | 955 +++++++++++++++++++++++++++++++++----- backend/go-app/walkoff.go | 10 +- 4 files changed, 853 insertions(+), 121 deletions(-) diff --git a/backend/go-app/go.mod b/backend/go-app/go.mod index 05925447..92e341a4 100644 --- a/backend/go-app/go.mod +++ b/backend/go-app/go.mod @@ -20,6 +20,9 @@ require ( github.com/docker/docker v20.10.3-0.20210216175712-646072ed6524+incompatible github.com/docker/go-connections v0.4.0 github.com/docker/go-units v0.4.0 // indirect + github.com/elastic/go-elasticsearch v0.0.0 // indirect + 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.40 github.com/fsouza/go-dockerclient v1.7.2 // indirect diff --git a/backend/go-app/go.sum b/backend/go-app/go.sum index 1cb3f682..b56cb465 100644 --- a/backend/go-app/go.sum +++ b/backend/go-app/go.sum @@ -109,6 +109,12 @@ github.com/docker/go-connections v0.4.0/go.mod h1:Gbd7IOopHjR8Iph03tsViu4nIes5Xh github.com/docker/go-units v0.4.0 h1:3uh0PgVws3nIA0Q+MwDC8yjEPf9zjRfZZWXZYDct3Tw= github.com/docker/go-units v0.4.0/go.mod h1:fgPhTUdO+D/Jk86RDLlptpiXQzgHJF7gydDDbaIK4Dk= github.com/dustin/go-humanize v1.0.0/go.mod h1:HtrtbFcZ19U5GC7JDqmcUSB87Iq5E25KnS6fMYU6eOk= +github.com/elastic/go-elasticsearch v0.0.0 h1:Pd5fqOuBxKxv83b0+xOAJDAkziWYwFinWnBO0y+TZaA= +github.com/elastic/go-elasticsearch v0.0.0/go.mod h1:TkBSJBuTyFdBnrNqoPc54FN0vKf5c04IdM4zuStJ7xg= +github.com/elastic/go-elasticsearch/v7 v7.12.0 h1:j4tvcMrZJLp39L2NYvBb7f+lHKPqPHSL3nvB8+/DV+s= +github.com/elastic/go-elasticsearch/v7 v7.12.0/go.mod h1:OJ4wdbtDNk5g503kvlHLyErCgQwwzmDtaFC4XyOxXA4= +github.com/elastic/go-elasticsearch/v8 v8.0.0-20210519083322-55daf7425ecb h1:svC8T5+v+aWpWiTt3nsGfpdqVb4NIWK/WamGXXECBXA= +github.com/elastic/go-elasticsearch/v8 v8.0.0-20210519083322-55daf7425ecb/go.mod h1:xe9a/L2aeOgFKKgrO3ibQTnMdpAeL0GC+5/HpGScSa4= github.com/emirpasic/gods v1.12.0 h1:QAUIPSaCu4G+POclxeqb3F+WPpdKqFGlw36+yOzGlrg= github.com/emirpasic/gods v1.12.0/go.mod h1:YfzfFFoVP/catgzJb4IKIqXjX78Ha8FMSDh3ymbK86o= github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4= diff --git a/backend/go-app/main.go b/backend/go-app/main.go index c961e799..81381ba3 100644 --- a/backend/go-app/main.go +++ b/backend/go-app/main.go @@ -35,9 +35,13 @@ import ( "cloud.google.com/go/storage" "google.golang.org/appengine/mail" + "github.com/elastic/go-elasticsearch/v8" + //"github.com/elastic/go-elasticsearch/v8/esapi" + "github.com/frikky/kin-openapi/openapi2" "github.com/frikky/kin-openapi/openapi2conv" "github.com/frikky/kin-openapi/openapi3" + /* "github.com/frikky/kin-openapi/openapi2" "github.com/frikky/kin-openapi/openapi2conv" @@ -639,14 +643,22 @@ func createNewUser(username, password, role, apikey string, org shuffle.OrgMini) } ctx := context.Background() - q := datastore.NewQuery("Users").Filter("Username =", username) - var users []shuffle.User - _, err = dbclient.GetAll(ctx, q, &users) + //users, err := FindUser(ctx context.Context, username string) ([]User, error) { + + users, err := shuffle.FindUser(ctx, strings.ToLower(strings.TrimSpace(username))) if err != nil && len(users) == 0 { - log.Printf("[WARNING] Failed getting user for registration: %s", err) + log.Printf("[WARNING] Failed getting user %s: %s", username, err) return err } + //q := datastore.NewQuery("Users").Filter("Username =", username) + //var users []shuffle.User + //_, err = dbclient.GetAll(ctx, q, &users) + //if err != nil && len(users) == 0 { + // log.Printf("[WARNING] Failed getting user for registration: %s", err) + // return err + //} + if len(users) > 0 { return errors.New(fmt.Sprintf("Username %s already exists", username)) } @@ -742,7 +754,9 @@ func handleRegister(resp http.ResponseWriter, request *http.Request) { // FIXME: Overhaul the top part. // Only admin can CREATE users, but if there are no users, anyone can make (first) - count, countErr := shuffle.GetUserCount() + ctx := context.Background() + users, countErr := shuffle.GetAllUsers(ctx) + count := len(users) user, err := shuffle.HandleApiAuthentication(resp, request) if err != nil { if (countErr == nil && count > 0) || countErr != nil { @@ -774,27 +788,23 @@ func handleRegister(resp http.ResponseWriter, request *http.Request) { role = "admin" } - ctx := context.Background() currentOrg := user.ActiveOrg if user.ActiveOrg.Id == "" { - log.Printf("There's no active org for the user. Checking if there's a single one to assing it to.") + log.Printf("[WARNING] There's no active org for the user %s. Checking if there's a single one to assing it to.", user.Username) - var orgs []shuffle.Org - q := datastore.NewQuery("Organizations") - _, err = dbclient.GetAll(ctx, q, &orgs) - if err == nil && len(orgs) == 1 { - log.Printf("No org exists in auth. Setting to default (first one)") + orgs, err := shuffle.GetAllOrgs(ctx) + if err == nil && len(orgs) > 0 { + log.Printf("[WARNING] No org exists for user %s. Setting to default (first one)", user.Username) currentOrg = shuffle.OrgMini{ Id: orgs[0].Id, Name: orgs[0].Name, } } - } err = createNewUser(data.Username, data.Password, role, "", currentOrg) if err != nil { - log.Printf("Failed registering user: %s", err) + log.Printf("[WARNING] Failed registering user: %s", err) resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return @@ -802,7 +812,7 @@ func handleRegister(resp http.ResponseWriter, request *http.Request) { resp.WriteHeader(200) resp.Write([]byte(`{"success": true}`)) - log.Printf("%s Successfully registered.", data.Username) + log.Printf("[INFO] %s Successfully registered.", data.Username) } func handleCookie(request *http.Request) bool { @@ -833,13 +843,6 @@ func handleInfo(resp http.ResponseWriter, request *http.Request) { } ctx := context.Background() - //session, err := getSession(ctx, userInfo.Session) - //if err != nil { - // log.Printf("Session %#v doesn't exist: %s", session, err) - // resp.WriteHeader(401) - // resp.Write([]byte(`{"success": false, "reason": "No session"}`)) - // return - //} // This is a long check to see if an inactive admin can access the site parsedAdmin := "false" @@ -852,9 +855,10 @@ func handleInfo(resp http.ResponseWriter, request *http.Request) { parsedAdmin = "true" ctx := context.Background() - q := datastore.NewQuery("Users") - var users []shuffle.User - _, err = dbclient.GetAll(ctx, q, &users) + users, err := shuffle.GetAllUsers(ctx) + //q := datastore.NewQuery("Users") + //var users []shuffle.User + //_, err = dbclient.GetAll(ctx, q, &users) if err != nil { resp.WriteHeader(401) resp.Write([]byte(`{"success": false, "reason": "Failed to get other users when verifying admin user"}`)) @@ -920,9 +924,10 @@ func handleInfo(resp http.ResponseWriter, request *http.Request) { if (len(userInfo.ActiveOrg.Name) == 0 || len(userInfo.ActiveOrg.Id) == 0) && len(userInfo.Orgs) > 0 { _, err := shuffle.GetOrg(ctx, userInfo.Orgs[0]) if err != nil { - var orgs []shuffle.Org - q := datastore.NewQuery("Organizations") - _, err = dbclient.GetAll(ctx, q, &orgs) + orgs, err := shuffle.GetAllOrgs(ctx) + //var orgs []shuffle.Org + //q := datastore.NewQuery("Organizations") + //_, err = dbclient.GetAll(ctx, q, &orgs) if err == nil { newStringOrgs := []string{} newOrgs := []shuffle.Org{} @@ -1056,13 +1061,16 @@ func checkAdminLogin(resp http.ResponseWriter, request *http.Request) { return } - count, err := shuffle.GetUserCount() + ctx := context.Background() + users, err := shuffle.GetAllUsers(ctx) if err != nil { resp.WriteHeader(401) resp.Write([]byte(fmt.Sprintf(`{"success": false, "reason": "%s"}`, err))) return } + count := len(users) + if count == 0 { log.Printf("[WARNING] No users - redirecting for management user") resp.WriteHeader(200) @@ -1099,16 +1107,24 @@ func handleLogin(resp http.ResponseWriter, request *http.Request) { ctx := context.Background() log.Printf("[INFO] Login Username: %s", data.Username) - q := datastore.NewQuery("Users").Filter("Username =", data.Username) - var users []shuffle.User - _, err = dbclient.GetAll(ctx, q, &users) - if err != nil { - log.Printf("Failed getting user %s", data.Username) + users, err := shuffle.FindUser(ctx, strings.ToLower(strings.TrimSpace(data.Username))) + if err != nil && len(users) == 0 { + log.Printf("[WARNING] Failed getting user %s: %s", data.Username, err) resp.WriteHeader(401) resp.Write([]byte(`{"success": false, "reason": "Username and/or password is incorrect"}`)) return } + //q := datastore.NewQuery("Users").Filter("Username =", data.Username) + //var users []shuffle.User + //_, err = dbclient.GetAll(ctx, q, &users) + //if err != nil { + // log.Printf("Failed getting user %s", data.Username) + // resp.WriteHeader(401) + // resp.Write([]byte(`{"success": false, "reason": "Username and/or password is incorrect"}`)) + // return + //} + if len(users) != 1 { log.Printf(`Found multiple or no users with the same username: %s: %d`, data.Username, len(users)) resp.WriteHeader(401) @@ -1191,39 +1207,6 @@ func handleLogin(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(loginData)) } -func setOpenApiDatastore(ctx context.Context, id string, data ParsedOpenApi) error { - k := datastore.NameKey("openapi3", id, nil) - if _, err := dbclient.Put(ctx, k, &data); err != nil { - return err - } - - return nil -} - -func getOpenApiDatastore(ctx context.Context, id string) (ParsedOpenApi, error) { - key := datastore.NameKey("openapi3", id, nil) - api := &ParsedOpenApi{} - if err := dbclient.Get(ctx, key, api); err != nil { - return ParsedOpenApi{}, err - } - - return *api, nil -} - -func setEnvironment(ctx context.Context, data *shuffle.Environment) error { - // clear session_token and API_token for user - k := datastore.NameKey("Environments", strings.ToLower(data.Name), nil) - - // New struct, to not add body, author etc - - if _, err := dbclient.Put(ctx, k, data); err != nil { - log.Println(err) - return err - } - - return nil -} - func fixOrgUser(ctx context.Context, org *shuffle.Org) *shuffle.Org { //found := false //for _, id := range user.Orgs { @@ -1719,7 +1702,7 @@ func setSpecificSchedule(resp http.ResponseWriter, request *http.Request) { } jsonPrettyPrint(string(body)) - var schedule ScheduleOld + var schedule shuffle.ScheduleOld err = json.Unmarshal(body, &schedule) if err != nil { log.Printf("Failed unmarshaling: %s", err) @@ -1730,7 +1713,7 @@ func setSpecificSchedule(resp http.ResponseWriter, request *http.Request) { // FIXME - check access etc ctx := context.Background() - err = setSchedule(ctx, schedule) + err = shuffle.SetSchedule(ctx, schedule) if err != nil { log.Printf("Failed setting schedule: %s", err) resp.WriteHeader(401) @@ -1875,9 +1858,9 @@ func handleNewSchedule(resp http.ResponseWriter, request *http.Request) { // FIXME - timestamp! // FIXME - applocation - cloud function? timeNow := int64(time.Now().Unix()) - schedule := ScheduleOld{ + schedule := shuffle.ScheduleOld{ Id: newId, - AppInfo: AppInfo{}, + AppInfo: shuffle.AppInfo{}, BaseAppLocation: "/home/frikky/git/shaffuru/tmp/apps", CreationTime: timeNow, LastModificationtime: timeNow, @@ -1885,7 +1868,7 @@ func handleNewSchedule(resp http.ResponseWriter, request *http.Request) { } ctx := context.Background() - err := setSchedule(ctx, schedule) + err := shuffle.SetSchedule(ctx, schedule) if err != nil { log.Printf("Failed setting hook: %s", err) resp.WriteHeader(401) @@ -2476,19 +2459,6 @@ func uploadWorkflowResult(resp http.ResponseWriter, request *http.Request) { resp.Write([]byte(`{"success": true}`)) } -// Index = Username -func setSchedule(ctx context.Context, schedule ScheduleOld) error { - key1 := datastore.NameKey("schedules", strings.ToLower(schedule.Id), nil) - - // New struct, to not add body, author etc - if _, err := dbclient.Put(ctx, key1, &schedule); err != nil { - log.Printf("Error adding schedule: %s", err) - return err - } - - return nil -} - //dst: {name: "title", required: "true", type: "string"} // //"title": "symptomDescription", @@ -3188,7 +3158,7 @@ func getOpenapi(resp http.ResponseWriter, request *http.Request) { // log.Println("You're supposed to be able to continue now.") //} - parsedApi, err := getOpenApiDatastore(ctx, id) + parsedApi, err := shuffle.GetOpenApiDatastore(ctx, id) if err != nil { resp.WriteHeader(401) resp.Write([]byte(`{"success": false}`)) @@ -3278,7 +3248,7 @@ func echoOpenapiData(resp http.ResponseWriter, request *http.Request) { resp.Write(urlbody) } -func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { +func handleSwaggerValidation(body []byte) (shuffle.ParsedOpenApi, error) { type versionCheck struct { Swagger string `datastore:"swagger" json:"swagger" yaml:"swagger"` SwaggerVersion string `datastore:"swaggerVersion" json:"swaggerVersion" yaml:"swaggerVersion"` @@ -3299,7 +3269,7 @@ func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { // support map[string]interface and similar (openapi3.Swagger) var version versionCheck - parsed := ParsedOpenApi{} + parsed := shuffle.ParsedOpenApi{} swaggerdata := []byte{} idstring := "" @@ -3329,13 +3299,13 @@ func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { swaggerv3, err := swaggerLoader.LoadSwaggerFromData(body) if err != nil { log.Printf("Failed parsing OpenAPI: %s", err) - return ParsedOpenApi{}, err + return shuffle.ParsedOpenApi{}, err } swaggerdata, err = json.Marshal(swaggerv3) if err != nil { log.Printf("Failed unmarshaling v3 data: %s", err) - return ParsedOpenApi{}, err + return shuffle.ParsedOpenApi{}, err } hasher := md5.New() @@ -3353,7 +3323,7 @@ func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { err = yaml.Unmarshal(body, &swagger) if err != nil { log.Printf("Yaml error (2): %s", err) - return ParsedOpenApi{}, err + return shuffle.ParsedOpenApi{}, err } else { //log.Printf("Valid yaml!") } @@ -3363,13 +3333,13 @@ func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { swaggerv3, err := openapi2conv.ToV3Swagger(&swagger) if err != nil { log.Printf("Failed converting from openapi2 to 3: %s", err) - return ParsedOpenApi{}, err + return shuffle.ParsedOpenApi{}, err } swaggerdata, err = json.Marshal(swaggerv3) if err != nil { log.Printf("Failed unmarshaling v3 data: %s", err) - return ParsedOpenApi{}, err + return shuffle.ParsedOpenApi{}, err } hasher := md5.New() @@ -3386,7 +3356,7 @@ func handleSwaggerValidation(body []byte) (ParsedOpenApi, error) { body = swaggerdata // Parsing it to swagger 3 - parsed = ParsedOpenApi{ + parsed = shuffle.ParsedOpenApi{ ID: idstring, Body: string(body), Success: true, @@ -3681,7 +3651,7 @@ func verifySwagger(resp http.ResponseWriter, request *http.Request) { } //log.Printf("DO I REACH HERE WHEN SAVING?") - parsed := ParsedOpenApi{ + parsed := shuffle.ParsedOpenApi{ ID: newmd5, Body: string(body), } @@ -3690,7 +3660,7 @@ func verifySwagger(resp http.ResponseWriter, request *http.Request) { // FIXME: Might cause versioning issues if we re-use the same!! // FIXME: Need a way to track different versions of the same app properly. // Hint: Save API.id somewhere, and use newmd5 to save latest version - err = setOpenApiDatastore(ctx, newmd5, parsed) + err = shuffle.SetOpenApiDatastore(ctx, newmd5, parsed) if err != nil { log.Printf("[ERROR] Failed saving to datastore: %s", err) resp.WriteHeader(500) @@ -3698,7 +3668,7 @@ func verifySwagger(resp http.ResponseWriter, request *http.Request) { } // Backup every single one - setOpenApiDatastore(ctx, api.ID, parsed) + shuffle.SetOpenApiDatastore(ctx, api.ID, parsed) /* err = increaseStatisticsField(ctx, "total_apps_created", newmd5, 1, user.ActiveOrg.Id) @@ -4152,16 +4122,650 @@ func remoteOrgJobHandler(org shuffle.Org, interval int) error { return nil } +func runInitEs(ctx context.Context) { + log.Printf("[DEBUG] Starting INIT setup (ES)") + httpProxy := os.Getenv("HTTP_PROXY") + if len(httpProxy) > 0 { + log.Printf("Running with HTTP proxy %s (env: HTTP_PROXY)", httpProxy) + } + httpsProxy := os.Getenv("HTTPS_PROXY") + if len(httpsProxy) > 0 { + log.Printf("Running with HTTPS proxy %s (env: HTTPS_PROXY)", httpsProxy) + } + + setUsers := false + log.Printf("[DEBUG] Getting organizations") + activeOrgs, err := shuffle.GetAllOrgs(ctx) + log.Printf("ORGS: %d", len(activeOrgs)) + if err != nil { + log.Printf("Error getting organizations: %s", err) + } else { + // Add all users to it + if len(activeOrgs) == 1 { + setUsers = true + } + + if len(activeOrgs) == 0 { + log.Printf(`No orgs. Setting 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, + } + + 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 + } + } else { + log.Printf("[DEBUG] There are %d org(s).", len(activeOrgs)) + + if len(activeOrgs) == 1 { + if len(activeOrgs[0].Users) == 0 { + log.Printf("ORG doesn't have any users??") + + q := datastore.NewQuery("Users") + var users []shuffle.User + _, err = dbclient.GetAll(ctx, q, &users) + if err != nil && len(users) == 0 { + log.Printf("Failed getting users in org fix") + } else { + // Remapping everyone to admin. This should never happen. + + for _, user := range users { + user.ActiveOrg = shuffle.OrgMini{ + Id: activeOrgs[0].Id, + Name: activeOrgs[0].Name, + Role: "admin", + } + + activeOrgs[0].Users = append(activeOrgs[0].Users, user) + } + + err = shuffle.SetOrg(ctx, activeOrgs[0], activeOrgs[0].Id) + if err != nil { + log.Printf("Failed setting org: %s", err) + } else { + log.Printf("Successfully updated org to have users!") + } + } + + } + } + } + } + _ = setUsers + + // Adding the users to the base organization since only one exists (default) + //.if setUsers && len(activeOrgs) > 0 { + //. activeOrg := activeOrgs[0] + + //. q := datastore.NewQuery("Users") + //. var users []shuffle.User + //. _, err = dbclient.GetAll(ctx, q, &users) + //. if err == nil { + //. setOrgBool := false + //. usernames := []string{} + //. for _, user := range users { + //. usernames = append(usernames, user.Username) + //. newUser := shuffle.User{ + //. Username: user.Username, + //. Id: user.Id, + //. ActiveOrg: shuffle.OrgMini{ + //. Id: activeOrg.Id, + //. }, + //. Orgs: []string{activeOrg.Id}, + //. Role: user.Role, + //. } + + //. found := false + //. for _, orgUser := range activeOrg.Users { + //. if user.Id == orgUser.Id { + //. found = true + //. } + //. } + + //. if !found && len(user.Username) > 0 { + //. log.Printf("Adding user %s to org %s", user.Username, activeOrg.Name) + //. activeOrg.Users = append(activeOrg.Users, newUser) + //. setOrgBool = true + //. } + //. } + + //. log.Printf("Users found: %s", strings.Join(usernames, ", ")) + + //. if setOrgBool { + //. err = shuffle.SetOrg(ctx, activeOrg, activeOrg.Id) + //. if err != nil { + //. log.Printf("Failed setting org %s: %s!", activeOrg.Name, err) + //. } else { + //. log.Printf("UPDATED org %s!", activeOrg.Name) + //. } + //. } + //. } + + //. log.Printf("Should add %d users to organization default", len(users)) + //.} + + //.if len(activeOrgs) == 0 { + //. orgQuery := datastore.NewQuery("Organizations") + //. _, err = dbclient.GetAll(ctx, orgQuery, &activeOrgs) + //. if err != nil { + //. log.Printf("Failed getting orgs the second time around") + //. } + //.} + + //.// Fix active users etc + //.q := datastore.NewQuery("Users").Filter("active =", true) + //.var activeusers []shuffle.User + //._, err = dbclient.GetAll(ctx, q, &activeusers) + //.if err != nil && len(activeusers) == 0 { + //. log.Printf("Error getting users during init: %s", err) + //.} else { + //. log.Printf("Parsing all users and setting them to active.") + //. q := datastore.NewQuery("Users") + //. var users []shuffle.User + //. _, err := dbclient.GetAll(ctx, q, &users) + //. //log.Printf("User ret: %s", err) + + //. if len(activeusers) == 0 && len(users) > 0 { + //. log.Printf("No active users found - setting ALL to active") + //. if err == nil { + //. for _, user := range users { + //. user.Active = true + //. if len(user.Username) == 0 { + //. shuffle.DeleteKey(ctx, "Users", strings.ToLower(user.Username)) + //. continue + //. } + + //. if len(user.Role) > 0 { + //. user.Roles = append(user.Roles, user.Role) + //. } + + //. if len(user.Orgs) == 0 { + //. defaultName := "default" + //. user.Orgs = []string{defaultName} + //. user.ActiveOrg = shuffle.OrgMini{ + //. Name: defaultName, + //. Role: "admin", + //. } + //. } + + //. err = shuffle.SetUser(ctx, &user, true) + //. if err != nil { + //. log.Printf("Failed to reset user") + //. } else { + //. log.Printf("Remade user %s with ID", user.Id) + //. err = shuffle.DeleteKey(ctx, "Users", strings.ToLower(user.Username)) + //. if err != nil { + //. log.Printf("Failed to delete old user by username") + //. } + //. } + //. } + //. } + //. } else if len(users) == 0 { + //. log.Printf("Trying to set up user based on environments SHUFFLE_DEFAULT_USERNAME & SHUFFLE_DEFAULT_PASSWORD") + //. username := os.Getenv("SHUFFLE_DEFAULT_USERNAME") + //. password := os.Getenv("SHUFFLE_DEFAULT_PASSWORD") + //. if len(username) == 0 || len(password) == 0 { + //. log.Printf("SHUFFLE_DEFAULT_USERNAME and SHUFFLE_DEFAULT_PASSWORD not defined as environments. Running without default user.") + //. } else { + //. apikey := os.Getenv("SHUFFLE_DEFAULT_APIKEY") + + //. tmpOrg := shuffle.OrgMini{ + //. Name: "default", + //. } + + //. err = createNewUser(username, password, "admin", apikey, tmpOrg) + //. if err != nil { + //. log.Printf("Failed to create default user %s: %s", username, err) + //. } else { + //. log.Printf("Successfully created user %s", username) + //. } + //. } + //. } else { + //. if len(users) < 10 && len(users) > 0 { + //. for _, user := range users { + //. log.Printf("[INFO] Username: %s, role: %s", user.Username, user.Role) + //. } + //. } else { + //. log.Printf("Found %d users.", len(users)) + //. } + + //. if len(activeOrgs) == 1 && len(users) > 0 { + //. for _, user := range users { + //. if user.ActiveOrg.Id == "" && len(user.Username) > 0 { + //. user.ActiveOrg = shuffle.OrgMini{ + //. Id: activeOrgs[0].Id, + //. Name: activeOrgs[0].Name, + //. } + + //. err = shuffle.SetUser(ctx, &user, true) + //. if err != nil { + //. log.Printf("Failed updating user %s with org", user.Username) + //. } else { + //. log.Printf("Updated user %s to have org", user.Username) + //. } + //. } + //. } + //. } + //. //log.Printf(users[0].Username) + //. } + //.} + + //.// Gets environments and inits if it doesn't exist + //.count, err := shuffle.GetEnvironmentCount() + //.if count == 0 && err == nil && len(activeOrgs) == 1 { + //. log.Printf("[INFO] Setting up environment with org %s", activeOrgs[0].Id) + //. item := shuffle.Environment{ + //. Name: "Shuffle", + //. Type: "onprem", + //. OrgId: activeOrgs[0].Id, + //. Default: true, + //. Id: uuid.NewV4().String(), + //. } + + //. err = setEnvironment(ctx, &item) + //. if err != nil { + //. log.Printf("Failed setting up new environment") + //. } + //.} else if len(activeOrgs) == 1 { + //. log.Printf("[INFO] Setting up all environments with org %s", activeOrgs[0].Id) + //. var environments []shuffle.Environment + //. q := datastore.NewQuery("Environments") + //. _, err = dbclient.GetAll(ctx, q, &environments) + //. if err == nil { + //. existingEnv := []string{} + //. _ = existingEnv + //. for _, item := range environments { + //. //if shuffle.ArrayContains(existingEnv, item.Name) { + //. // log.Printf("[WARNING] Env %s already exists - deleting it. %#v", item.Name, item) + //. // err = DeleteKey(ctx, "Environments", item.Name) + //. // if err != nil { + //. // log.Printf("[WARNING] Env deletion error: %s", err) + //. // } + + //. // continue + //. //} + + //. //existingEnv = append(existingEnv, item.Name) + + //. if item.OrgId == activeOrgs[0].Id && len(item.Id) > 0 { + //. continue + //. } + + //. if len(item.Id) == 0 { + //. item.Id = uuid.NewV4().String() + //. } + + //. item.OrgId = activeOrgs[0].Id + //. err = setEnvironment(ctx, &item) + //. if err != nil { + //. log.Printf("Failed adding environment to org %s", activeOrgs[0].Id) + //. } + //. } + //. } + //.} + + //.// Fixing workflows to have real activeorg IDs + //.//workflowQ := datastore.NewQuery("workflow") + //.//ret, err := dbclient.GetAll(ctx, workflowQ, &workflows) + //.//log.Printf("[INFO] Found %d workflows during startup", workflowCount) + //.//log.Printf("%#v, %s", ret, err) + + //.var workflows []shuffle.Workflow + //.if len(activeOrgs) == 1 { + //. q := datastore.NewQuery("workflow").Limit(35) + //. _, err = dbclient.GetAll(ctx, q, &workflows) + //. if err != nil && len(workflows) == 0 { + //. log.Printf("Error getting workflows in runinit: %s", err) + //. } else { + //. updated := 0 + //. timeNow := time.Now().Unix() + //. for _, workflow := range workflows { + //. setLocal := false + //. if workflow.ExecutingOrg.Id == "" || len(workflow.OrgId) == 0 { + //. workflow.OrgId = activeOrgs[0].Id + //. workflow.ExecutingOrg = shuffle.OrgMini{ + //. Id: activeOrgs[0].Id, + //. Name: activeOrgs[0].Name, + //. } + + //. setLocal = true + //. } else if workflow.Edited == 0 { + //. workflow.Edited = timeNow + //. setLocal = true + //. } + + //. if setLocal { + //. err = shuffle.SetWorkflow(ctx, workflow, workflow.ID) + //. if err != nil { + //. log.Printf("Failed setting workflow in init: %s", err) + //. } else { + //. log.Printf("Fixed workflow %s to have the right info.", workflow.ID) + //. updated += 1 + //. } + //. } + //. } + + //. if updated > 0 { + //. log.Printf("Set workflow orgs for %d workflows", updated) + //. } + //. } + + //. /* + //. fileq := datastore.NewQuery("Files").Limit(1) + //. count, err := dbclient.Count(ctx, fileq) + //. log.Printf("FILECOUNT: %d", count) + //. if err == nil && count < 10 { + //. basepath := "." + //. filename := "testfile.txt" + //. fileId := uuid.NewV4().String() + //. log.Printf("Creating new file reference %s because none exist!", fileId) + //. workflowId := "2cf1169d-b460-41de-8c36-28b2092866f8" + //. downloadPath := fmt.Sprintf("%s/%s/%s/%s", basepath, activeOrgs[0].Id, workflowId, fileId) + + //. timeNow := time.Now().Unix() + //. newFile := File{ + //. Id: fileId, + //. CreatedAt: timeNow, + //. UpdatedAt: timeNow, + //. Description: "Created by system for testing", + //. Status: "active", + //. Filename: filename, + //. OrgId: activeOrgs[0].Id, + //. WorkflowId: workflowId, + //. DownloadPath: downloadPath, + //. } + + //. err = setFile(ctx, newFile) + //. if err != nil { + //. log.Printf("Failed setting file: %s", err) + //. } else { + //. log.Printf("Created file %s in init", newFile.DownloadPath) + //. } + //. } + //. */ + + //. var allworkflowapps []shuffle.AppAuthenticationStorage + //. q = datastore.NewQuery("workflowappauth") + //. _, err = dbclient.GetAll(ctx, q, &allworkflowapps) + //. if err == nil { + //. log.Printf("Setting up all app auths with org %s", activeOrgs[0].Id) + //. for _, item := range allworkflowapps { + //. if item.OrgId != "" { + //. continue + //. } + + //. //log.Printf("Should update auth for %#v!", item) + //. item.OrgId = activeOrgs[0].Id + //. err = shuffle.SetWorkflowAppAuthDatastore(ctx, item, item.Id) + //. if err != nil { + //. log.Printf("Failed adding AUTH to org %s", activeOrgs[0].Id) + //. } + //. } + //. } + + //. var schedules []ScheduleOld + //. q = datastore.NewQuery("schedules") + //. _, err = dbclient.GetAll(ctx, q, &schedules) + //. if err == nil { + //. log.Printf("Setting up all schedules with org %s", activeOrgs[0].Id) + //. for _, item := range schedules { + //. if item.Org != "" { + //. continue + //. } + + //. log.Printf("ENV: %s", item.Environment) + //. if item.Environment == "cloud" { + //. log.Printf("Skipping cloud schedule") + //. continue + //. } + + //. item.Org = activeOrgs[0].Id + //. err = setSchedule(ctx, item) + //. if err != nil { + //. log.Printf("Failed adding schedule to org %s", activeOrgs[0].Id) + //. } + //. } + //. } + //.} + + //.log.Printf("Starting cloud schedules for orgs!") + //.type requestStruct struct { + //. ApiKey string `json:"api_key"` + //.} + //.for _, org := range activeOrgs { + //. if !org.CloudSync { + //. log.Printf("Skipping org %s because sync isn't set (1).", org.Id) + //. continue + //. } + + //. //interval := int(org.SyncConfig.Interval) + //. interval := 15 + //. if interval == 0 { + //. log.Printf("Skipping org %s because sync isn't set (0).", org.Id) + //. continue + //. } + + //. log.Printf("Should start schedule for org %s", org.Name) + //. job := func() { + //. err := remoteOrgJobHandler(org, interval) + //. if err != nil { + //. log.Printf("[ERROR] Failed request with remote org setup (2): %s", err) + //. } + //. } + + //. jobret, err := newscheduler.Every(int(interval)).Seconds().NotImmediately().Run(job) + //. if err != nil { + //. log.Printf("[CRITICAL] Failed to schedule org: %s", err) + //. } else { + //. log.Printf("Started sync on interval %d for org %s", interval, org.Name) + //. scheduledOrgs[org.Id] = jobret + //. } + //.} + + //.// Gets schedules and starts them + //.log.Printf("Relaunching schedules") + //.schedules, err := shuffle.GetAllSchedules(ctx, "ALL") + //.if err != nil { + //. log.Printf("Failed getting schedules during service init: %s", err) + //.} else { + //. log.Printf("Setting up %d schedule(s)", len(schedules)) + //. url := &url.URL{} + //. for _, schedule := range schedules { + //. if schedule.Environment == "cloud" { + //. log.Printf("Skipping cloud schedule") + //. continue + //. } + + //. //log.Printf("Schedule: %#v", schedule) + //. job := func() { + //. //log.Printf("[INFO] Running schedule %s with interval %d.", schedule.Id, schedule.Seconds) + //. //log.Printf("ARG: %s", schedule.WrappedArgument) + + //. request := &http.Request{ + //. URL: url, + //. Method: "POST", + //. Body: ioutil.NopCloser(strings.NewReader(schedule.WrappedArgument)), + //. } + + //. _, _, err := handleExecution(schedule.WorkflowId, shuffle.Workflow{}, request) + //. if err != nil { + //. log.Printf("[WARNING] Failed to execute %s: %s", schedule.WorkflowId, err) + //. } + //. } + + //. //log.Printf("Schedule time: every %d seconds", schedule.Seconds) + //. jobret, err := newscheduler.Every(schedule.Seconds).Seconds().NotImmediately().Run(job) + //. if err != nil { + //. log.Printf("Failed to schedule workflow: %s", err) + //. } + + //. scheduledJobs[schedule.Id] = jobret + //. } + //.} + + //.// form force-flag to download workflow apps + //.forceUpdateEnv := os.Getenv("SHUFFLE_APP_FORCE_UPDATE") + //.forceUpdate := false + //.if len(forceUpdateEnv) > 0 && forceUpdateEnv == "true" { + //. log.Printf("Forcing to rebuild apps") + //. forceUpdate = true + //.} + + //.// Getting apps to see if we should initialize a test + //.log.Printf("[INFO] Getting and validating workflowapps") + //.workflowapps, err := shuffle.GetAllWorkflowApps(ctx, 500) + //.if err != nil && len(workflowapps) == 0 { + //. log.Printf("[WARNING] Failed getting apps (runInit): %s", err) + //.} else if err == nil && len(workflowapps) > 0 { + //. var allworkflowapps []shuffle.WorkflowApp + //. q := datastore.NewQuery("workflowapp") + //. _, err := dbclient.GetAll(ctx, q, &allworkflowapps) + //. if err == nil { + //. for _, workflowapp := range allworkflowapps { + //. if workflowapp.Edited == 0 { + //. err = shuffle.SetWorkflowAppDatastore(ctx, workflowapp, workflowapp.ID) + //. if err == nil { + //. log.Printf("[INFO] Updating time for workflowapp %s:%s", workflowapp.Name, workflowapp.AppVersion) + //. } + //. } + //. } + //. } + + //.} else if err == nil && len(workflowapps) == 0 { + //. log.Printf("Downloading default workflow apps") + //. fs := memfs.New() + //. storer := memory.NewStorage() + + //. url := os.Getenv("SHUFFLE_APP_DOWNLOAD_LOCATION") + //. if len(url) == 0 { + //. url = "https://github.com/frikky/shuffle-apps" + //. } + + //. username := os.Getenv("SHUFFLE_DOWNLOAD_AUTH_USERNAME") + //. password := os.Getenv("SHUFFLE_DOWNLOAD_AUTH_PASSWORD") + + //. cloneOptions := &git.CloneOptions{ + //. URL: url, + //. } + + //. if len(username) > 0 && len(password) > 0 { + //. cloneOptions.Auth = &http2.BasicAuth{ + //. Username: username, + //. Password: password, + //. } + //. } + //. branch := os.Getenv("SHUFFLE_DOWNLOAD_AUTH_BRANCH") + //. if len(branch) > 0 && branch != "master" && branch != "main" { + //. cloneOptions.ReferenceName = plumbing.ReferenceName(branch) + //. } + + //. log.Printf("Getting apps from %s", url) + + //. r, err := git.Clone(storer, fs, cloneOptions) + + //. if err != nil { + //. log.Printf("Failed loading repo into memory (init): %s", err) + //. } + + //. dir, err := fs.ReadDir("") + //. if err != nil { + //. log.Printf("Failed reading folder: %s", err) + //. } + //. _ = r + //. //iterateAppGithubFolders(fs, dir, "", "testing") + + //. // FIXME: Get all the apps? + //. _, _, err = iterateAppGithubFolders(fs, dir, "", "", forceUpdate) + //. if err != nil { + //. log.Printf("[WARNING] Error from app load in init: %s", err) + //. } + //. //_, _, err = iterateAppGithubFolders(fs, dir, "", "", forceUpdate) + + //. // Hotloads locally + //. location := os.Getenv("SHUFFLE_APP_HOTLOAD_FOLDER") + //. if len(location) != 0 { + //. handleAppHotload(ctx, location, false) + //. } + //.} + + //.log.Printf("[INFO] Downloading OpenAPI data for search - EXTRA APPS") + //.apis := "https://github.com/frikky/security-openapis" + + //.// THis gets memory problems hahah + //.//apis := "https://github.com/APIs-guru/openapi-directory" + //.fs := memfs.New() + //.storer := memory.NewStorage() + //.cloneOptions := &git.CloneOptions{ + //. URL: apis, + //.} + //._, err = git.Clone(storer, fs, cloneOptions) + //.if err != nil { + //. log.Printf("Failed loading repo %s into memory: %s", apis, err) + //.} else { + //. log.Printf("[INFO] Finished git clone. Looking for updates to the repo.") + //. dir, err := fs.ReadDir("") + //. if err != nil { + //. log.Printf("Failed reading folder: %s", err) + //. } + + //. iterateOpenApiGithub(fs, dir, "", "") + //. log.Printf("[INFO] Finished downloading extra API samples") + //.} + + //.workflowLocation := os.Getenv("SHUFFLE_DOWNLOAD_WORKFLOW_LOCATION") + //.if len(workflowLocation) > 0 { + //. log.Printf("[INFO] Downloading WORKFLOWS from %s if no workflows - EXTRA workflows", workflowLocation) + //. q := datastore.NewQuery("workflow").Limit(35) + //. var workflows []shuffle.Workflow + //. _, err = dbclient.GetAll(ctx, q, &workflows) + //. if err != nil && len(workflows) == 0 { + //. log.Printf("Error getting workflows: %s", err) + //. } else { + //. if len(workflows) == 0 { + //. username := os.Getenv("SHUFFLE_DOWNLOAD_WORKFLOW_USERNAME") + //. password := os.Getenv("SHUFFLE_DOWNLOAD_WORKFLOW_PASSWORD") + //. orgId := "" + //. if len(activeOrgs) > 0 { + //. orgId = activeOrgs[0].Id + //. } + + //. err = loadGithubWorkflows(workflowLocation, username, password, "", os.Getenv("SHUFFLE_DOWNLOAD_WORKFLOW_BRANCH"), orgId) + //. if err != nil { + //. log.Printf("Failed to upload workflows from github: %s", err) + //. } else { + //. log.Printf("[INFO] Finished downloading workflows from github!") + //. } + //. } else { + //. log.Printf("[INFO] Skipping because there are %d workflows already", len(workflows)) + //. } + + //. } + //.} + + log.Printf("[INFO] Finished INIT") +} + // Handles configuration items during Shuffle startup func runInit(ctx context.Context) { // Setting stats for backend starts (failure count as well) - log.Printf("[DEBUG] Starting INIT setup") - err := increaseStatisticsField(ctx, "backend_executions", "", 1, "") - if err != nil { - log.Printf("Failed increasing local stats: %s", err) - } - log.Printf("[DEBUG] Finalized init statistics update") + //err := increaseStatisticsField(ctx, "backend_executions", "", 1, "") + //if err != nil { + // log.Printf("Failed increasing local stats: %s", err) + //} + //log.Printf("[DEBUG] Finalized init statistics update") + log.Printf("[DEBUG] Starting INIT setup") httpProxy := os.Getenv("HTTP_PROXY") if len(httpProxy) > 0 { log.Printf("Running with HTTP proxy %s (env: HTTP_PROXY)", httpProxy) @@ -4210,7 +4814,7 @@ func runInit(ctx context.Context) { log.Printf("[DEBUG] Getting organizations") orgQuery := datastore.NewQuery("Organizations") var activeOrgs []shuffle.Org - _, err = dbclient.GetAll(ctx, orgQuery, &activeOrgs) + _, err := dbclient.GetAll(ctx, orgQuery, &activeOrgs) if err != nil { log.Printf("Error getting organizations!") } else { @@ -4241,7 +4845,7 @@ func runInit(ctx context.Context) { setUsers = true } } else { - log.Printf("There are %d org(s).", len(activeOrgs)) + log.Printf("[DEBUG] There are %d org(s).", len(activeOrgs)) if len(activeOrgs) == 1 { if len(activeOrgs[0].Users) == 0 { @@ -4447,9 +5051,9 @@ func runInit(ctx context.Context) { Id: uuid.NewV4().String(), } - err = setEnvironment(ctx, &item) + err = shuffle.SetEnvironment(ctx, &item) if err != nil { - log.Printf("Failed setting up new environment") + log.Printf("[WARNING] Failed setting up new environment") } } else if len(activeOrgs) == 1 { log.Printf("[INFO] Setting up all environments with org %s", activeOrgs[0].Id) @@ -4481,9 +5085,9 @@ func runInit(ctx context.Context) { } item.OrgId = activeOrgs[0].Id - err = setEnvironment(ctx, &item) + err = shuffle.SetEnvironment(ctx, &item) if err != nil { - log.Printf("Failed adding environment to org %s", activeOrgs[0].Id) + log.Printf("[WARNING] Failed adding environment to org %s", activeOrgs[0].Id) } } } @@ -4588,7 +5192,7 @@ func runInit(ctx context.Context) { } } - var schedules []ScheduleOld + var schedules []shuffle.ScheduleOld q = datastore.NewQuery("schedules") _, err = dbclient.GetAll(ctx, q, &schedules) if err == nil { @@ -4605,7 +5209,7 @@ func runInit(ctx context.Context) { } item.Org = activeOrgs[0].Id - err = setSchedule(ctx, item) + err = shuffle.SetSchedule(ctx, item) if err != nil { log.Printf("Failed adding schedule to org %s", activeOrgs[0].Id) } @@ -4944,7 +5548,7 @@ func handleStopCloudSync(syncUrl string, org shuffle.Org) (*shuffle.Org, error) if environment.Type == "cloud" { environment.Name = "Cloud" environment.Archived = true - err = setEnvironment(ctx, &environment) + err = shuffle.SetEnvironment(ctx, &environment) if err == nil { log.Printf("[INFO] Updated cloud environment %s", environment.Name) } else { @@ -5197,7 +5801,7 @@ func handleCloudSetup(resp http.ResponseWriter, request *http.Request) { if environment.Type == "cloud" { environment.Name = "Cloud" environment.Archived = false - err = setEnvironment(ctx, &environment) + err = shuffle.SetEnvironment(ctx, &environment) if err == nil { log.Printf("[INFO] Re-added cloud environment %s", environment.Name) } else { @@ -5221,7 +5825,7 @@ func handleCloudSetup(resp http.ResponseWriter, request *http.Request) { Id: uuid.NewV4().String(), } - err = setEnvironment(ctx, &newEnv) + err = shuffle.SetEnvironment(ctx, &newEnv) if err != nil { log.Printf("Failed setting up NEW org environment for org %s: %s", org.Id, err) } else { @@ -5250,6 +5854,114 @@ func handleCloudSetup(resp http.ResponseWriter, request *http.Request) { resp.Write(respBody) } +func migrateDatabase(resp http.ResponseWriter, request *http.Request) { + cors := shuffle.HandleCors(resp, request) + if cors { + return + } + + user, userErr := shuffle.HandleApiAuthentication(resp, request) + if userErr != nil { + log.Printf("[WARNING] Api authentication failed in make workflow public: %s", userErr) + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + if user.Role != "admin" { + log.Printf("[WARNING] Failed to migrate because you're not admin") + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + + type dbSetup struct { + dbUrl string + } + + //setup := dbSetup{ + // dbUrl: "http://192.168.3.8:9200", + //} + + config := elasticsearch.Config{ + Addresses: []string{ + "https://192.168.3.8:9200", + }, + Username: "", + Password: "", + } + + //es, err := elasticsearch.NewDefaultClient() + es, err := elasticsearch.NewClient(config) + if err != nil { + log.Printf("[WARNING] Failed connecting with es7: %s", err) + resp.WriteHeader(401) + resp.Write([]byte(`{"success": false}`)) + return + } + _ = es + + /* + q := datastore.NewQuery(indexType) + var users []shuffle.User + _, err = dbclient.GetAll(ctx, q, &users) + if err == nil && len(users) > 0 { + } else { + log.Printf("[WARNING] Failed getting USERS: %s", err) + + for _, user := range users { + id := user.Id + log.Printf("Indexing %s (%s)", user.Username, id) + + b, err := json.Marshal(user) + if err != nil { + log.Printf("[WARNING] Failed marshalling %s - %s: %s", id, indexType, err) + return + } + + req := esapi.IndexRequest{ + Index: indexType, + DocumentID: id, + Body: strings.NewReader(string(b)), + Refresh: "true", + } + + res, err := req.Do(context.Background(), es) + if err != nil { + log.Printf("Error getting response: %s", err) + } + + defer res.Body.Close() + if res.IsError() { + log.Printf("[%s] Error indexing document ID=%d", res.Status(), id) + } else { + log.Printf("Successfully indexed %s of ID %s", indexType, id) + } + + break + } + } + */ + + // 1. Get users + // 2. Get organizations + // 3. Get files + // 4. Get workflows + // 5. Get apps + // 6. Get workflowappauth + // 7. Get workflowexecution + // 9. Get Environments + // 10. Get hooks + // 11. Get openapi3 + // 12. Get schedules + // 13. Get sessions + // 14. Get workflowqueue + + //log.Printf("[INFO] Successfully published workflow %s (%s) TO CLOUD", workflow.Name, workflow.ID) + resp.WriteHeader(200) + resp.Write([]byte(fmt.Sprintf(`{"success": true}`))) +} + func makeWorkflowPublic(resp http.ResponseWriter, request *http.Request) { cors := shuffle.HandleCors(resp, request) if cors { @@ -5401,10 +6113,20 @@ func initHandlers() { panic(fmt.Sprintf("[DEBUG] Database client error during init: %s", err)) } - _ = shuffle.RunInit(*dbclient, storage.Client{}, gceProject, "onprem", true) + es, err := elasticsearch.NewClient( + elasticsearch.Config{ + Addresses: []string{"http://127.0.0.1:9200"}, + }, + ) + if err != nil { + panic(fmt.Sprintf("[DEBUG] Database client for ELASTICSEARCH error during init: %s", err)) + } + + _ = shuffle.RunInit(*dbclient, *es, storage.Client{}, gceProject, "onprem", true, "elasticsearch") log.Printf("[DEBUG] Finished Shuffle database init") - go runInit(ctx) + //go runInit(ctx) + go runInitEs(ctx) r := mux.NewRouter() r.HandleFunc("/api/v1/_ah/health", healthCheckHandler) @@ -5522,6 +6244,7 @@ func initHandlers() { // Docker orborus specific r.HandleFunc("/api/v1/get_docker_image", getDockerImage).Methods("POST", "OPTIONS") + r.HandleFunc("/api/v1/migrate_database", migrateDatabase).Methods("POST", "OPTIONS") // Important for email, IDS etc. Create this by: // PS: For cloud, this has to use cloud storage. diff --git a/backend/go-app/walkoff.go b/backend/go-app/walkoff.go index 64cec153..506d47ad 100644 --- a/backend/go-app/walkoff.go +++ b/backend/go-app/walkoff.go @@ -660,7 +660,7 @@ func createSchedule(ctx context.Context, scheduleId, workflowId, name, startNode // Doesn't need running/not running. If stopped, we just delete it. timeNow := int64(time.Now().Unix()) - schedule := ScheduleOld{ + schedule := shuffle.ScheduleOld{ Id: scheduleId, WorkflowId: workflowId, StartNode: startNode, @@ -674,7 +674,7 @@ func createSchedule(ctx context.Context, scheduleId, workflowId, name, startNode Environment: "onprem", } - err = setSchedule(ctx, schedule) + err = shuffle.SetSchedule(ctx, schedule) if err != nil { log.Printf("Failed to set schedule: %s", err) return err @@ -2828,7 +2828,7 @@ func scheduleWorkflow(resp http.ResponseWriter, request *http.Request) { } timeNow := int64(time.Now().Unix()) - newSchedule := ScheduleOld{ + newSchedule := shuffle.ScheduleOld{ Id: schedule.Id, WorkflowId: workflow.ID, StartNode: startNode, @@ -2842,7 +2842,7 @@ func scheduleWorkflow(resp http.ResponseWriter, request *http.Request) { Environment: "cloud", } - err = setSchedule(ctx, newSchedule) + err = shuffle.SetSchedule(ctx, newSchedule) if err != nil { log.Printf("Failed setting cloud schedule: %s", err) resp.WriteHeader(401) @@ -3883,7 +3883,7 @@ func iterateOpenApiGithub(fs billy.Filesystem, dir []os.FileInfo, extra string, log.Printf("Added %s:%s to the database from OpenAPI repo", api.Name, api.AppVersion) // Set OpenAPI datastore - err = setOpenApiDatastore(ctx, parsedOpenApi.ID, parsedOpenApi) + err = shuffle.SetOpenApiDatastore(ctx, parsedOpenApi.ID, parsedOpenApi) if err != nil { log.Printf("Failed uploading openapi to datastore in loop: %s", err) continue