diff --git a/cmd/compose/main.go b/cmd/compose/main.go index a054d8f65..ccbf992a8 100644 --- a/cmd/compose/main.go +++ b/cmd/compose/main.go @@ -2,11 +2,18 @@ package main import ( "github.com/cortezaproject/corteza-server/compose" - "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/corteza" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/logger" ) func main() { - cfg := compose.Configure() - cmd := cfg.MakeCLI(cli.Context()) - cli.HandleError(cmd.Execute()) + logger.Init() + + app.Run( + logger.Default(), + app.NewOptions(compose.SERVICE), + &corteza.App{}, + &compose.App{}, + ) } diff --git a/cmd/messaging/main.go b/cmd/messaging/main.go index 35578dfe8..4412b0d79 100644 --- a/cmd/messaging/main.go +++ b/cmd/messaging/main.go @@ -1,12 +1,19 @@ package main import ( + "github.com/cortezaproject/corteza-server/corteza" "github.com/cortezaproject/corteza-server/messaging" - "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/logger" ) func main() { - cfg := messaging.Configure() - cmd := cfg.MakeCLI(cli.Context()) - cli.HandleError(cmd.Execute()) + logger.Init() + + app.Run( + logger.Default(), + app.NewOptions(messaging.SERVICE), + &corteza.App{}, + &messaging.App{}, + ) } diff --git a/cmd/monolith/main.go b/cmd/monolith/main.go index 3ff8d945a..c0759c11e 100644 --- a/cmd/monolith/main.go +++ b/cmd/monolith/main.go @@ -1,12 +1,26 @@ package main import ( + "github.com/cortezaproject/corteza-server/compose" + "github.com/cortezaproject/corteza-server/corteza" + "github.com/cortezaproject/corteza-server/messaging" "github.com/cortezaproject/corteza-server/monolith" - "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/logger" + "github.com/cortezaproject/corteza-server/system" ) func main() { - cfg := monolith.Configure() - cmd := cfg.MakeCLI(cli.Context()) - cli.HandleError(cmd.Execute()) + logger.Init() + + app.Run( + logger.Default(), + app.NewOptions(), + &corteza.App{}, + &monolith.App{ + System: &system.App{}, + Compose: &compose.App{}, + Messaging: &messaging.App{}, + }, + ) } diff --git a/cmd/system/main.go b/cmd/system/main.go index 1f7d5f5f8..cb9bce2f3 100644 --- a/cmd/system/main.go +++ b/cmd/system/main.go @@ -1,12 +1,19 @@ package main import ( - "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/corteza" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/logger" "github.com/cortezaproject/corteza-server/system" ) func main() { - cfg := system.Configure() - cmd := cfg.MakeCLI(cli.Context()) - cli.HandleError(cmd.Execute()) + logger.Init() + + app.Run( + logger.Default(), + app.NewOptions(system.SERVICE), + &corteza.App{}, + &system.App{}, + ) } diff --git a/codegen.sh b/codegen.sh index a1759920c..86040012b 100755 --- a/codegen.sh +++ b/codegen.sh @@ -99,12 +99,6 @@ function types { ./build/gen-type-set-test --types Rule --output pkg/permissions/rule.gen_test.go --with-primary-key=false --package permissions ./build/gen-type-set-test --types Resource --output pkg/permissions/resource.gen_test.go --with-primary-key=false --package permissions - ./build/gen-type-set --types Script --output pkg/automation/script.gen.go --package automation - ./build/gen-type-set-test --types Script --output pkg/automation/script.gen_test.go --package automation - ./build/gen-type-set --types Trigger --output pkg/automation/trigger.gen.go --package automation - ./build/gen-type-set-test --types Trigger --output pkg/automation/trigger.gen_test.go --package automation - - green "OK" } diff --git a/compose/app.go b/compose/app.go new file mode 100644 index 000000000..139a58b4d --- /dev/null +++ b/compose/app.go @@ -0,0 +1,96 @@ +package compose + +import ( + "context" + + "github.com/go-chi/chi" + _ "github.com/joho/godotenv/autoload" + "github.com/spf13/cobra" + "github.com/titpetric/factory" + "go.uber.org/zap" + + "github.com/cortezaproject/corteza-server/compose/commands" + migrate "github.com/cortezaproject/corteza-server/compose/db" + "github.com/cortezaproject/corteza-server/compose/rest" + "github.com/cortezaproject/corteza-server/compose/service" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/system/auth/external" +) + +type ( + App struct { + Opts *app.Options + Log *zap.Logger + } +) + +const SERVICE = "compose" + +func (app *App) Setup(log *zap.Logger, opts *app.Options) (err error) { + app.Log = log.Named(SERVICE) + app.Opts = opts + return +} + +func (app *App) Upgrade(ctx context.Context) (err error) { + db := factory.Database.MustGet(SERVICE, "default").With(ctx).Quiet() + err = migrate.Migrate(db, app.Log) + if err != nil { + return + } + + return +} + +// Initialized +func (app *App) Initialize(ctx context.Context) (err error) { + // Connects to all services it needs to + err = service.Initialize(ctx, app.Log, service.Config{ + Storage: app.Opts.Storage, + Corredor: app.Opts.Corredor, + }) + + if err != nil { + return + } + + // Initialize external authentication (from default settings) + external.Init() + return +} + +func (app *App) Activate(ctx context.Context) (err error) { + if err = service.Activate(ctx); err != nil { + return + } + + service.Watchers(ctx) + + // Wire in cross service JWT maker for Corredor + corredor.Service().SetJwtMaker(corredor.CrossServiceAuthTokenMaker()) + + return +} + +func (app *App) Provision(ctx context.Context) (err error) { + ctx = auth.SetSuperUserContext(ctx) + + if err = provisionConfig(ctx, app.Log); err != nil { + return + } + + return +} + +func (app *App) MountApiRoutes(r chi.Router) { + rest.MountRoutes(r) +} + +func (app *App) RegisterCliCommands(p *cobra.Command) { + p.AddCommand( + commands.Importer(), + commands.Exporter(), + ) +} diff --git a/compose/app_test.go b/compose/app_test.go new file mode 100644 index 000000000..8de4be6e2 --- /dev/null +++ b/compose/app_test.go @@ -0,0 +1,11 @@ +package compose + +import ( + "testing" + + "github.com/cortezaproject/corteza-server/pkg/app" +) + +func TestConfigure(t *testing.T) { + var _ app.Runnable = &App{} +} diff --git a/compose/commands/exporter.go b/compose/commands/exporter.go index 5d9aae89b..6fa58dabb 100644 --- a/compose/commands/exporter.go +++ b/compose/commands/exporter.go @@ -28,19 +28,15 @@ import ( sysTypes "github.com/cortezaproject/corteza-server/system/types" ) -func Exporter(ctx context.Context, c *cli.Config) *cobra.Command { +func Exporter() *cobra.Command { cmd := &cobra.Command{ Use: "export", Short: "Export", Long: `Specify one ("modules", "pages", "charts", "permissions") or more resources to export`, Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - - ctx = auth.SetSuperUserContext(ctx) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) nsFlag = cmd.Flags().Lookup("namespace").Value.String() sFlag = cmd.Flags().Lookup("settings").Changed pFlag = cmd.Flags().Lookup("permissions").Changed @@ -117,19 +113,9 @@ func nsExporter(ctx context.Context, out *Compose, nsFlag string, args []string) charts, _, err := service.DefaultChart.Find(types.ChartFilter{NamespaceID: ns.ID}) cli.HandleError(err) - scripts, _, err := service.DefaultInternalAutomationManager.FindScripts(ctx, automation.ScriptFilter{}) - cli.HandleError(err) - - triggers, _, err := service.DefaultInternalAutomationManager.FindTriggers(ctx, automation.TriggerFilter{}) - cli.HandleError(err) - - scripts, _ = scripts.Filter(func(script *automation.Script) (b bool, e error) { - return script.NamespaceID == ns.ID, nil - }) - // nsOut.Name = ns.Name // nsOut.Handle = ns.Slug - // nsOut.Enabled = ns.Enabled + // nsOut.Always = ns.Always // nsOut.Meta = ns.Meta // // nsOut.Allow = sysExporter.ExportableResourcePermissions(roles, service.DefaultPermissions, permissions.Allow, ns.PermissionResource()) @@ -142,9 +128,7 @@ func nsExporter(ctx context.Context, out *Compose, nsFlag string, args []string) case "chart", "charts": nsOut.Charts = expCharts(charts, modules) case "page", "pages": - nsOut.Pages = expPages(0, pages, modules, charts, scripts) - case "scripts", "triggers", "automation": - nsOut.Scripts = expAutomation(scripts, triggers, modules) + nsOut.Pages = expPages(0, pages, modules, charts) } } @@ -419,7 +403,7 @@ func expModuleFieldOptions(f *types.ModuleField, modules types.ModuleSet) types. return out } -func expPages(parentID uint64, pages types.PageSet, modules types.ModuleSet, charts types.ChartSet, scripts automation.ScriptSet) (o yaml.MapSlice) { +func expPages(parentID uint64, pages types.PageSet, modules types.ModuleSet, charts types.ChartSet) (o yaml.MapSlice) { var ( children = pages.FindByParent(parentID) handle string @@ -430,8 +414,8 @@ func expPages(parentID uint64, pages types.PageSet, modules types.ModuleSet, cha page := Page{ Title: child.Title, Description: child.Description, - Blocks: expPageBlocks(child.Blocks, pages, modules, charts, scripts), - Pages: expPages(child.ID, pages, modules, charts, scripts), + Blocks: expPageBlocks(child.Blocks, pages, modules, charts), + Pages: expPages(child.ID, pages, modules, charts), Visible: child.Visible, Allow: sysExporter.ExportableResourcePermissions(roles, service.DefaultPermissions, permissions.Allow, types.PagePermissionResource), @@ -464,7 +448,7 @@ func expPages(parentID uint64, pages types.PageSet, modules types.ModuleSet, cha return } -func expPageBlocks(in types.PageBlocks, pages types.PageSet, modules types.ModuleSet, charts types.ChartSet, scripts automation.ScriptSet) types.PageBlocks { +func expPageBlocks(in types.PageBlocks, pages types.PageSet, modules types.ModuleSet, charts types.ChartSet) types.PageBlocks { out := types.PageBlocks(in) // Remove extra options to keep the output tidy @@ -536,9 +520,9 @@ func expPageBlocks(in types.PageBlocks, pages types.PageSet, modules types.Modul _ = deinterfacer.Each(btn, func(_ int, k string, v interface{}) error { switch k { case "triggerID", "scriptID": - if s := scripts.FindByID(deinterfacer.ToUint64(v)); s != nil { - button["script"] = makeHandleFromName(s.Name, "", "automation-script-%d", s.ID) - } + // if s := scripts.FindByID(deinterfacer.ToUint64(v)); s != nil { + // button["script"] = makeHandleFromName(s.Name, "", "automation-script-%d", s.ID) + // } default: button[k] = v } diff --git a/compose/commands/importer.go b/compose/commands/importer.go index d09f784f5..baf8f04f7 100644 --- a/compose/commands/importer.go +++ b/compose/commands/importer.go @@ -1,7 +1,6 @@ package commands import ( - "context" "io" "os" "strconv" @@ -16,16 +15,14 @@ import ( "github.com/cortezaproject/corteza-server/pkg/cli" ) -func Importer(ctx context.Context, c *cli.Config) *cobra.Command { +func Importer() *cobra.Command { cmd := &cobra.Command{ Use: "import", Short: "Import", Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) ff []io.Reader nsFlag = cmd.Flags().Lookup("namespace").Value.String() ns *types.Namespace @@ -45,8 +42,6 @@ func Importer(ctx context.Context, c *cli.Config) *cobra.Command { } } - ctx = auth.SetSuperUserContext(ctx) - if len(args) > 0 { ff = make([]io.Reader, len(args)) for a, arg := range args { diff --git a/compose/compose.go b/compose/compose.go deleted file mode 100644 index 4ee57a4b7..000000000 --- a/compose/compose.go +++ /dev/null @@ -1,93 +0,0 @@ -package compose - -import ( - "context" - - _ "github.com/joho/godotenv/autoload" - "github.com/spf13/cobra" - "github.com/titpetric/factory" - - "github.com/cortezaproject/corteza-server/compose/commands" - migrate "github.com/cortezaproject/corteza-server/compose/db" - "github.com/cortezaproject/corteza-server/compose/rest" - "github.com/cortezaproject/corteza-server/compose/service" - "github.com/cortezaproject/corteza-server/pkg/cli" -) - -const ( - compose = "compose" -) - -func Configure() *cli.Config { - var servicesInitialized bool - - return &cli.Config{ - ServiceName: compose, - - RootCommandPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - return - }, - }, - - InitServices: func(ctx context.Context, c *cli.Config) { - if servicesInitialized { - return - } - servicesInitialized = true - - cli.HandleError(service.Init(ctx, c.Log, service.Config{ - Storage: *c.StorageOpt, - Corredor: *c.ScriptRunner, - GRPCClientSystem: *c.GRPCServerSystem, - })) - }, - - ApiServerPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - if c.ProvisionOpt.MigrateDatabase { - cli.HandleError(c.ProvisionMigrateDatabase.Run(ctx, cmd, c)) - } - - c.InitServices(ctx, c) - - if c.ProvisionOpt.Configuration { - cli.HandleError(provisionConfig(ctx, cmd, c)) - } - - go service.Watchers(ctx) - return nil - }, - }, - - ApiServerRoutes: cli.Mounters{ - rest.MountRoutes, - }, - - AdtSubCommands: cli.CommandMakers{ - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Importer(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Exporter(ctx, c) - }, - }, - - ProvisionMigrateDatabase: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - var db, err = factory.Database.Get(compose) - if err != nil { - return err - } - - db = db.With(ctx).Quiet() - - return migrate.Migrate(db, c.Log) - }, - }, - - ProvisionConfig: cli.Runners{ - provisionConfig, - }, - } -} diff --git a/compose/compose_test.go b/compose/compose_test.go deleted file mode 100644 index cea01b871..000000000 --- a/compose/compose_test.go +++ /dev/null @@ -1,15 +0,0 @@ -package compose - -import ( - "context" - "testing" - - "github.com/stretchr/testify/require" -) - -func TestConfigure(t *testing.T) { - var config = Configure() - require.True(t, config != nil, "Configure valid") - require.True(t, func() bool { config.Init(); return true }(), "Initialization ok") - require.True(t, config.MakeCLI(context.Background()) != nil, "CLI created") -} diff --git a/compose/provision-config.go b/compose/provision.go similarity index 85% rename from compose/provision-config.go rename to compose/provision.go index f9d0c0040..5daa08abc 100644 --- a/compose/provision-config.go +++ b/compose/provision.go @@ -5,22 +5,20 @@ import ( "io" "github.com/pkg/errors" - "github.com/spf13/cobra" + "go.uber.org/zap" "gopkg.in/yaml.v2" "github.com/cortezaproject/corteza-server/compose/importer" "github.com/cortezaproject/corteza-server/compose/service" "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli" impAux "github.com/cortezaproject/corteza-server/pkg/importer" "github.com/cortezaproject/corteza-server/pkg/settings" provision "github.com/cortezaproject/corteza-server/provision/compose" ) -func provisionConfig(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - c.Log.Debug("running configuration provision") - c.InitServices(ctx, c) +func provisionConfig(ctx context.Context, log *zap.Logger) (err error) { + log.Debug("running configuration provision") var provisioned bool @@ -30,7 +28,7 @@ func provisionConfig(ctx context.Context, cmd *cobra.Command, c *cli.Config) (er if provisioned, err = isProvisioned(ctx); err != nil { return err } else if provisioned { - c.Log.Debug("configuration already provisioned") + log.Debug("configuration already provisioned") } readers, err := impAux.ReadStatic(provision.Asset) @@ -64,7 +62,7 @@ func partialImportSettings(ctx context.Context, ss settings.Service, ff ...io.Re // importer w/o permissions & roles // we need only settings - imp = importer.NewImporter(nil, nil, nil, nil, nil, nil, si) + imp = importer.NewImporter(nil, nil, nil, nil, nil, si) // current value current settings.ValueSet diff --git a/compose/repository/repository.go b/compose/repository/repository.go index 2765b5adf..306cb1edb 100644 --- a/compose/repository/repository.go +++ b/compose/repository/repository.go @@ -15,7 +15,7 @@ type ( // DB produces a contextual DB handle func DB(ctx context.Context) *factory.DB { - return factory.Database.MustGet("compose").With(ctx) + return factory.Database.MustGet("compose", "default").With(ctx) } // With updates repository and database contexts diff --git a/compose/service/service.go b/compose/service/service.go index 38cbcfa27..a83969e0c 100644 --- a/compose/service/service.go +++ b/compose/service/service.go @@ -8,10 +8,10 @@ import ( "github.com/cortezaproject/corteza-server/compose/repository" "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/app/options" "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/pkg/automation" "github.com/cortezaproject/corteza-server/pkg/automation/corredor" - "github.com/cortezaproject/corteza-server/pkg/cli/options" "github.com/cortezaproject/corteza-server/pkg/permissions" "github.com/cortezaproject/corteza-server/pkg/settings" "github.com/cortezaproject/corteza-server/pkg/store" @@ -80,7 +80,8 @@ var ( DefaultSystemRole *systemRole ) -func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { +// Initializes compose-only services +func Initialize(ctx context.Context, log *zap.Logger, c Config) (err error) { var db = repository.DB(ctx) DefaultLogger = log.Named("service") @@ -90,6 +91,7 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { // to allow integration tests to inject own permission service DefaultPermissions = permissions.Service(ctx, DefaultLogger, db, "compose_permission_rules") } + DefaultAccessControl = AccessControl(DefaultPermissions) DefaultSettings = settings.NewService( @@ -99,12 +101,6 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { CurrentSettings, ) - // Run initial update of current settings with super-user credentials - err = DefaultSettings.UpdateCurrent(auth.SetSuperUserContext(ctx)) - if err != nil { - return - } - if DefaultStore == nil { if c.Storage.MinioEndpoint != "" { if c.Storage.MinioBucket == "" { @@ -202,6 +198,16 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { return nil } +func Activate(ctx context.Context) (err error) { + // Run initial update of current settings with super-user credentials + err = DefaultSettings.UpdateCurrent(auth.SetSuperUserContext(ctx)) + if err != nil { + return + } + + return +} + func Watchers(ctx context.Context) { // Reloading automation scripts on change DefaultAutomationRunner.Watch(ctx) diff --git a/corteza/corteza.go b/corteza/corteza.go new file mode 100644 index 000000000..0b984f08e --- /dev/null +++ b/corteza/corteza.go @@ -0,0 +1,83 @@ +package corteza + +import ( + "context" + "time" + + "github.com/pkg/errors" + "go.uber.org/zap" + + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/pkg/db" + "github.com/cortezaproject/corteza-server/pkg/http" + "github.com/cortezaproject/corteza-server/pkg/logger" + "github.com/cortezaproject/corteza-server/pkg/mail" + "github.com/cortezaproject/corteza-server/pkg/monitor" + "github.com/cortezaproject/corteza-server/pkg/sentry" +) + +type ( + App struct { + opt *app.Options + log *zap.Logger + } +) + +var _ app.Runnable = &App{} + +func (app *App) Setup(log *zap.Logger, opts *app.Options) (err error) { + app.log = log + app.opt = opts + + logger.SetDefault(log) + + if err = sentry.Init(opts.Sentry); err != nil { + return errors.Wrap(err, "could not initialize Sentry") + } + + // Use Sentry right away to handle any panics + // that might occur inside auth, mail setup... + defer sentry.Recover() + + auth.SetupDefault(opts.JWT.Secret, int(opts.JWT.Expiry/time.Minute)) + mail.SetupDialer(opts.SMTP.Host, opts.SMTP.Port, opts.SMTP.User, opts.SMTP.Pass, opts.SMTP.From) + + http.SetupDefaults( + opts.HTTPClient.HttpClientTimeout, + opts.HTTPClient.ClientTSLInsecure, + ) + + monitor.Setup(app.log, opts.Monitor) + + return +} + +func (app *App) Initialize(ctx context.Context) (err error) { + + defer sentry.Recover() + + _, err = db.TryToConnect(ctx, app.log, app.opt.DB) + if err != nil { + return errors.Wrap(err, "could not connect to database") + } + + return +} + +func (app *App) Upgrade(ctx context.Context) error { + return nil +} + +func (app *App) Activate(ctx context.Context) (err error) { + if err = corredor.Start(ctx, app.log, app.opt.Corredor); err != nil { + return err + } + + return nil +} + +func (app *App) Provision(ctx context.Context) error { + return nil +} diff --git a/messaging/app.go b/messaging/app.go new file mode 100644 index 000000000..3452674b7 --- /dev/null +++ b/messaging/app.go @@ -0,0 +1,104 @@ +package messaging + +import ( + "context" + + "github.com/go-chi/chi" + _ "github.com/joho/godotenv/autoload" + "github.com/spf13/cobra" + "github.com/titpetric/factory" + "go.uber.org/zap" + + "github.com/cortezaproject/corteza-server/messaging/commands" + migrate "github.com/cortezaproject/corteza-server/messaging/db" + "github.com/cortezaproject/corteza-server/messaging/rest" + "github.com/cortezaproject/corteza-server/messaging/service" + "github.com/cortezaproject/corteza-server/messaging/websocket" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/corredor" +) + +type ( + App struct { + Opts *app.Options + Log *zap.Logger + + ws *websocket.Websocket + } +) + +const SERVICE = "messaging" + +func (app *App) Setup(log *zap.Logger, opts *app.Options) (err error) { + app.Log = log.Named(SERVICE) + app.Opts = opts + + app.ws = websocket.New(&websocket.Config{ + Timeout: opts.Websocket.Timeout, + PingTimeout: opts.Websocket.PingTimeout, + PingPeriod: opts.Websocket.PingPeriod, + }) + + return +} + +func (app *App) Upgrade(ctx context.Context) (err error) { + db := factory.Database.MustGet(SERVICE, "default").With(ctx).Quiet() + err = migrate.Migrate(db, app.Log) + if err != nil { + return + } + + return +} + +func (app *App) Initialize(ctx context.Context) (err error) { + // Connects to all services it needs to + err = service.Initialize(ctx, app.Log, service.Config{ + Storage: app.Opts.Storage, + Corredor: app.Opts.Corredor, + }) + + if err != nil { + return + } + + return +} + +func (app *App) Activate(ctx context.Context) (err error) { + if err = service.Activate(ctx); err != nil { + return + } + + service.Watchers(ctx) + websocket.Watch(ctx) + + // Wire in cross service JWT maker for Corredor + corredor.Service().SetJwtMaker(corredor.CrossServiceAuthTokenMaker()) + + return +} + +func (app *App) Provision(ctx context.Context) (err error) { + ctx = auth.SetSuperUserContext(ctx) + + if err = provisionConfig(ctx, app.Log); err != nil { + return + } + + return +} + +func (app *App) MountApiRoutes(r chi.Router) { + rest.MountRoutes(r) + app.ws.ApiServerRoutes(r) +} + +func (app *App) RegisterCliCommands(p *cobra.Command) { + p.AddCommand( + commands.Importer(), + commands.Exporter(), + ) +} diff --git a/messaging/app_test.go b/messaging/app_test.go new file mode 100644 index 000000000..72621675a --- /dev/null +++ b/messaging/app_test.go @@ -0,0 +1,11 @@ +package messaging + +import ( + "testing" + + "github.com/cortezaproject/corteza-server/pkg/app" +) + +func TestConfigure(t *testing.T) { + var _ app.Runnable = &App{} +} diff --git a/messaging/commands/exporter.go b/messaging/commands/exporter.go index d8a1c0fb0..d517230b3 100644 --- a/messaging/commands/exporter.go +++ b/messaging/commands/exporter.go @@ -16,18 +16,15 @@ import ( sysTypes "github.com/cortezaproject/corteza-server/system/types" ) -func Exporter(ctx context.Context, c *cli.Config) *cobra.Command { +func Exporter() *cobra.Command { cmd := &cobra.Command{ Use: "export", Short: "Export", Long: `Export messaging resources`, Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) sFlag = cmd.Flags().Lookup("settings").Changed pFlag = cmd.Flags().Lookup("permissions").Changed diff --git a/messaging/commands/importer.go b/messaging/commands/importer.go index 85d3d304e..6eea8a39c 100644 --- a/messaging/commands/importer.go +++ b/messaging/commands/importer.go @@ -1,7 +1,6 @@ package commands import ( - "context" "io" "os" @@ -12,16 +11,14 @@ import ( "github.com/cortezaproject/corteza-server/pkg/cli" ) -func Importer(ctx context.Context, c *cli.Config) *cobra.Command { +func Importer() *cobra.Command { cmd := &cobra.Command{ Use: "import", Short: "Import", Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) ff []io.Reader err error ) diff --git a/messaging/messaging.go b/messaging/messaging.go deleted file mode 100644 index 72f81a02d..000000000 --- a/messaging/messaging.go +++ /dev/null @@ -1,109 +0,0 @@ -package messaging - -import ( - "context" - - "github.com/go-chi/chi" - _ "github.com/joho/godotenv/autoload" - "github.com/spf13/cobra" - "github.com/titpetric/factory" - - "github.com/cortezaproject/corteza-server/messaging/commands" - migrate "github.com/cortezaproject/corteza-server/messaging/db" - "github.com/cortezaproject/corteza-server/messaging/rest" - "github.com/cortezaproject/corteza-server/messaging/service" - "github.com/cortezaproject/corteza-server/messaging/websocket" - "github.com/cortezaproject/corteza-server/pkg/cli" - "github.com/cortezaproject/corteza-server/pkg/cli/options" -) - -const ( - messaging = "messaging" -) - -func Configure() *cli.Config { - var ( - servicesInitialized bool - - // Websocket handler - ws *websocket.Websocket - ) - - return &cli.Config{ - ServiceName: messaging, - - RootCommandPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - return - }, - }, - - InitServices: func(ctx context.Context, c *cli.Config) { - if servicesInitialized { - return - } - servicesInitialized = true - - cli.HandleError(service.Init(ctx, c.Log, service.Config{ - Storage: *c.StorageOpt, - })) - }, - - ApiServerPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - if c.ProvisionOpt.MigrateDatabase { - cli.HandleError(c.ProvisionMigrateDatabase.Run(ctx, cmd, c)) - } - - if c.ProvisionOpt.Configuration { - cli.HandleError(provisionConfig(ctx, cmd, c)) - } - - c.InitServices(ctx, c) - - var websocketOpt = options.Websocket(messaging) - - ws = websocket.Init(ctx, &websocket.Config{ - Timeout: websocketOpt.Timeout, - PingTimeout: websocketOpt.PingTimeout, - PingPeriod: websocketOpt.PingPeriod, - }) - - go service.Watchers(ctx) - return nil - }, - }, - - ApiServerRoutes: cli.Mounters{ - rest.MountRoutes, - // Wrap in func() to assure ws is set when mounted - func(r chi.Router) { ws.ApiServerRoutes(r) }, - }, - - AdtSubCommands: cli.CommandMakers{ - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Importer(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Exporter(ctx, c) - }, - }, - - ProvisionMigrateDatabase: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - var db, err = factory.Database.Get(messaging) - if err != nil { - return err - } - - db = db.With(ctx).Quiet() - - return migrate.Migrate(db, c.Log) - }, - }, - - ProvisionConfig: cli.Runners{ - provisionConfig, - }, - } -} diff --git a/messaging/messaging_test.go b/messaging/messaging_test.go deleted file mode 100644 index b9d0d39ba..000000000 --- a/messaging/messaging_test.go +++ /dev/null @@ -1,15 +0,0 @@ -package messaging - -import ( - "context" - "testing" - - "github.com/stretchr/testify/require" -) - -func TestConfigure(t *testing.T) { - var config = Configure() - require.True(t, config != nil, "Configure valid") - require.True(t, func() bool { config.Init(); return true }(), "Initialization ok") - require.True(t, config.MakeCLI(context.Background()) != nil, "CLI created") -} diff --git a/messaging/provision-config.go b/messaging/provision.go similarity index 88% rename from messaging/provision-config.go rename to messaging/provision.go index 774fbbef3..3b28739a1 100644 --- a/messaging/provision-config.go +++ b/messaging/provision.go @@ -5,22 +5,20 @@ import ( "io" "github.com/pkg/errors" - "github.com/spf13/cobra" + "go.uber.org/zap" "gopkg.in/yaml.v2" "github.com/cortezaproject/corteza-server/messaging/importer" "github.com/cortezaproject/corteza-server/messaging/service" "github.com/cortezaproject/corteza-server/messaging/types" "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli" impAux "github.com/cortezaproject/corteza-server/pkg/importer" "github.com/cortezaproject/corteza-server/pkg/settings" provision "github.com/cortezaproject/corteza-server/provision/messaging" ) -func provisionConfig(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - c.Log.Debug("running configuration provision") - c.InitServices(ctx, c) +func provisionConfig(ctx context.Context, log *zap.Logger) (err error) { + log.Debug("running configuration provision") var provisioned bool @@ -30,7 +28,7 @@ func provisionConfig(ctx context.Context, cmd *cobra.Command, c *cli.Config) (er if provisioned, err = isProvisioned(ctx); err != nil { return err } else if provisioned { - c.Log.Debug("configuration already provisioned") + log.Debug("configuration already provisioned") } readers, err := impAux.ReadStatic(provision.Asset) diff --git a/messaging/repository/repository.go b/messaging/repository/repository.go index 3bb62cd94..31283089f 100644 --- a/messaging/repository/repository.go +++ b/messaging/repository/repository.go @@ -15,7 +15,7 @@ type ( // DB produces a contextual DB handle func DB(ctx context.Context) *factory.DB { - return factory.Database.MustGet("messaging").With(ctx) + return factory.Database.MustGet("messaging", "default").With(ctx) } // With updates repository and database contexts diff --git a/messaging/service/service.go b/messaging/service/service.go index 92c7a52ac..5cb5e47fb 100644 --- a/messaging/service/service.go +++ b/messaging/service/service.go @@ -53,7 +53,7 @@ var ( DefaultWebhook WebhookService ) -func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { +func Initialize(ctx context.Context, log *zap.Logger, c Config) (err error) { DefaultLogger = log.Named("service") if DefaultPermissions == nil { @@ -71,12 +71,6 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { CurrentSettings, ) - // Run initial update of current settings with super-user credentials - err = DefaultSettings.UpdateCurrent(intAuth.SetSuperUserContext(ctx)) - if err != nil { - return - } - if DefaultStore == nil { if c.Storage.MinioEndpoint != "" { if c.Storage.MinioBucket == "" { @@ -127,6 +121,16 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { return nil } +func Activate(ctx context.Context) (err error) { + // Run initial update of current settings with super-user credentials + err = DefaultSettings.UpdateCurrent(intAuth.SetSuperUserContext(ctx)) + if err != nil { + return + } + + return +} + func Watchers(ctx context.Context) { DefaultPermissions.Watch(ctx) } diff --git a/messaging/websocket/websocket.go b/messaging/websocket/websocket.go index 9e068f119..df0056a9d 100644 --- a/messaging/websocket/websocket.go +++ b/messaging/websocket/websocket.go @@ -18,7 +18,7 @@ type ( } ) -func (Websocket) New(config *Config) *Websocket { +func New(config *Config) *Websocket { ws := &Websocket{ config: config, } diff --git a/monolith/app.go b/monolith/app.go new file mode 100644 index 000000000..245179cd9 --- /dev/null +++ b/monolith/app.go @@ -0,0 +1,137 @@ +package monolith + +import ( + "context" + + "github.com/go-chi/chi" + _ "github.com/joho/godotenv/autoload" + "github.com/spf13/cobra" + "go.uber.org/zap" + "google.golang.org/grpc" + + "github.com/cortezaproject/corteza-server/compose" + "github.com/cortezaproject/corteza-server/corteza" + "github.com/cortezaproject/corteza-server/messaging" + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/system" +) + +type ( + App struct { + Core *corteza.App + System *system.App + Compose *compose.App + Messaging *messaging.App + } +) + +func (monolith *App) Setup(log *zap.Logger, opts *app.Options) (err error) { + + // Make sure system behaves properly + // + // This will alter the auth settings provision procedure + system.IsMonolith = true + + return app.RunSetup( + log, + opts, + monolith.System, + monolith.Compose, + monolith.Messaging, + ) +} + +func (monolith *App) Upgrade(ctx context.Context) (err error) { + return app.RunUpgrade( + ctx, + monolith.System, + monolith.Compose, + monolith.Messaging, + ) +} + +func (monolith App) Initialize(ctx context.Context) (err error) { + return app.RunInitialize( + ctx, + monolith.System, + monolith.Compose, + monolith.Messaging, + ) +} + +func (monolith *App) Activate(ctx context.Context) (err error) { + err = app.RunActivate( + ctx, + monolith.System, + monolith.Compose, + monolith.Messaging, + ) + + if err != nil { + return + } + + // Override JWT maker for Corredor with internal + corredor.Service().SetJwtMaker(corredor.InternalAuthTokenMaker()) + + return +} + +func (monolith *App) Provision(ctx context.Context) (err error) { + return app.RunProvision( + ctx, + monolith.System, + monolith.Compose, + monolith.Messaging, + ) +} + +func (monolith *App) MountApiRoutes(r chi.Router) { + r.Route("/system", func(r chi.Router) { + monolith.System.MountApiRoutes(r) + }) + + r.Route("/compose", func(r chi.Router) { + monolith.Compose.MountApiRoutes(r) + }) + + r.Route("/messaging", func(r chi.Router) { + monolith.Messaging.MountApiRoutes(r) + }) +} + +func (monolith *App) RegisterGrpcServices(srv *grpc.Server) { + monolith.System.RegisterGrpcServices(srv) +} + +// RegisterCliCommands on monolith will wrapp all commands +func (monolith *App) RegisterCliCommands(rootCmd *cobra.Command) { + systemCmd := &cobra.Command{ + Use: "system", + Aliases: []string{"sys"}, + Short: "Commands from system service", + } + + composeCmd := &cobra.Command{ + Use: "compose", + Aliases: []string{"cmp"}, + Short: "Commands from messaging service", + } + + messagingCmd := &cobra.Command{ + Use: "messaging", + Aliases: []string{"msg"}, + Short: "Commands from compose service", + } + + monolith.System.RegisterCliCommands(systemCmd) + monolith.Compose.RegisterCliCommands(composeCmd) + monolith.Messaging.RegisterCliCommands(messagingCmd) + + rootCmd.AddCommand( + systemCmd, + composeCmd, + messagingCmd, + ) +} diff --git a/monolith/app_test.go b/monolith/app_test.go new file mode 100644 index 000000000..dc807de8e --- /dev/null +++ b/monolith/app_test.go @@ -0,0 +1,11 @@ +package monolith + +import ( + "testing" + + "github.com/cortezaproject/corteza-server/pkg/app" +) + +func TestConfigure(t *testing.T) { + var _ app.Runnable = &App{} +} diff --git a/monolith/monolith.go b/monolith/monolith.go deleted file mode 100644 index d61c813d6..000000000 --- a/monolith/monolith.go +++ /dev/null @@ -1,123 +0,0 @@ -package monolith - -import ( - "context" - - "github.com/go-chi/chi" - _ "github.com/joho/godotenv/autoload" - "github.com/spf13/cobra" - - "github.com/cortezaproject/corteza-server/compose" - "github.com/cortezaproject/corteza-server/messaging" - "github.com/cortezaproject/corteza-server/pkg/api" - "github.com/cortezaproject/corteza-server/pkg/cli" - "github.com/cortezaproject/corteza-server/system" -) - -func Configure() *cli.Config { - cmp := compose.Configure() - msg := messaging.Configure() - sys := system.Configure() - - cmp.Init() - msg.Init() - sys.Init() - - // Set API as a monolith build - api.Monolith = true - - // Combines all three services/apps and makes them run as one monolith app - return &cli.Config{ - ServiceName: "", - - InitServices: func(ctx context.Context, c *cli.Config) { - cmp.InitServices(ctx, cmp) - msg.InitServices(ctx, cmp) - sys.InitServices(ctx, cmp) - }, - - RootCommandDBSetup: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - cli.HandleError(cmp.RootCommandDBSetup.Run(ctx, cmd, cmp)) - cli.HandleError(msg.RootCommandDBSetup.Run(ctx, cmd, msg)) - cli.HandleError(sys.RootCommandDBSetup.Run(ctx, cmd, sys)) - return - }, - }, - - RootCommandName: "corteza-server", - RootCommandPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - cli.HandleError(cmp.RootCommandPreRun.Run(ctx, cmd, cmp)) - cli.HandleError(msg.RootCommandPreRun.Run(ctx, cmd, msg)) - cli.HandleError(sys.RootCommandPreRun.Run(ctx, cmd, sys)) - return - }, - }, - - ApiServerPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - cli.HandleError(cmp.ApiServerPreRun.Run(ctx, cmd, cmp)) - cli.HandleError(msg.ApiServerPreRun.Run(ctx, cmd, msg)) - cli.HandleError(sys.ApiServerPreRun.Run(ctx, cmd, sys)) - return - }, - }, - - ApiServerRoutes: cli.Mounters{ - func(r chi.Router) { - r.Route("/compose", cmp.ApiServerRoutes.MountRoutes) - r.Route("/messaging", msg.ApiServerRoutes.MountRoutes) - r.Route("/system", sys.ApiServerRoutes.MountRoutes) - }, - }, - - AdtSubCommands: cli.CommandMakers{ - func(ctx context.Context, c *cli.Config) *cobra.Command { - if cc := cmp.AdtSubCommands; len(cc) > 0 { - sub := &cobra.Command{Use: "compose", Short: "Commands from compose service"} - sub.AddCommand(cc.Make(ctx, c)...) - return sub - } - - return nil - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - if cc := msg.AdtSubCommands; len(cc) > 0 { - sub := &cobra.Command{Use: "messaging", Short: "Commands from messaging service"} - sub.AddCommand(cc.Make(ctx, c)...) - return sub - } - - return nil - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - if cc := sys.AdtSubCommands; len(cc) > 0 { - sub := &cobra.Command{Use: "system", Short: "Commands from system service"} - sub.AddCommand(cc.Make(ctx, c)...) - return sub - } - - return nil - }, - }, - - ProvisionMigrateDatabase: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - cli.HandleError(sys.ProvisionMigrateDatabase.Run(ctx, cmd, sys)) - cli.HandleError(cmp.ProvisionMigrateDatabase.Run(ctx, cmd, cmp)) - cli.HandleError(msg.ProvisionMigrateDatabase.Run(ctx, cmd, msg)) - return - }, - }, - - ProvisionConfig: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - cli.HandleError(sys.ProvisionConfig.Run(ctx, cmd, sys)) - cli.HandleError(cmp.ProvisionConfig.Run(ctx, cmd, cmp)) - cli.HandleError(msg.ProvisionConfig.Run(ctx, cmd, msg)) - return - }, - }, - } -} diff --git a/monolith/monolith_test.go b/monolith/monolith_test.go deleted file mode 100644 index bcbdc0c55..000000000 --- a/monolith/monolith_test.go +++ /dev/null @@ -1,15 +0,0 @@ -package monolith - -import ( - "context" - "testing" - - "github.com/stretchr/testify/require" -) - -func TestConfigure(t *testing.T) { - var config = Configure() - require.True(t, config != nil, "Configure valid") - require.True(t, func() bool { config.Init(); return true }(), "Initialization ok") - require.True(t, config.MakeCLI(context.Background()) != nil, "CLI created") -} diff --git a/pkg/api/debug.go b/pkg/api/debug.go index d1d51f1cb..333129d87 100644 --- a/pkg/api/debug.go +++ b/pkg/api/debug.go @@ -7,16 +7,10 @@ import ( "runtime" "github.com/go-chi/chi" - "github.com/go-chi/chi/middleware" ) -func Debug(r chi.Router) { - r.Mount("/debug", middleware.Profiler()) - DebugRoutes(r) -} - -func DebugRoutes(r chi.Router) { - r.Get("/debug/routes", func(w http.ResponseWriter, req *http.Request) { +func debugRoutes(r chi.Routes) http.HandlerFunc { + return func(w http.ResponseWriter, req *http.Request) { var printRoutes func(chi.Routes, string) printRoutes = func(r chi.Routes, pfix string) { @@ -33,5 +27,5 @@ func DebugRoutes(r chi.Router) { } printRoutes(r, "") - }) + } } diff --git a/pkg/api/metrics.go b/pkg/api/metrics.go index 7069c0256..e9e8a93ca 100644 --- a/pkg/api/metrics.go +++ b/pkg/api/metrics.go @@ -9,12 +9,12 @@ import ( "github.com/prometheus/client_golang/prometheus/promhttp" ) -// Middleware is the request logger that provides metrics to prometheus -func Middleware(name string) func(http.Handler) http.Handler { +// MetricsMiddleware is the request logger that provides metrics to prometheus +func metricsMiddleware(name string) func(http.Handler) http.Handler { return chiprometheus.NewMiddleware(name) } -func Mount(r chi.Router, username, password string) { +func metricsMount(r chi.Router, username, password string) { r.Group(func(r chi.Router) { r.Use(basicauth.New("Metrics", map[string][]string{ username: {password}, diff --git a/pkg/api/middleware.go b/pkg/api/middleware.go index dcacf51c3..73f8b6557 100644 --- a/pkg/api/middleware.go +++ b/pkg/api/middleware.go @@ -8,12 +8,12 @@ import ( "github.com/go-chi/chi/middleware" "go.uber.org/zap" - sentryhttp "github.com/getsentry/sentry-go/http" + "github.com/getsentry/sentry-go/http" "github.com/cortezaproject/corteza-server/pkg/logger" ) -func Base(log *zap.Logger) []func(http.Handler) http.Handler { +func BaseMiddleware(log *zap.Logger) []func(http.Handler) http.Handler { return []func(http.Handler) http.Handler{ handleCORS, middleware.RealIP, @@ -22,14 +22,14 @@ func Base(log *zap.Logger) []func(http.Handler) http.Handler { } } -func Sentry() func(http.Handler) http.Handler { +func sentryMiddleware() func(http.Handler) http.Handler { return sentryhttp.New(sentryhttp.Options{ Repanic: true, }).Handle } // HandlePanic sends 500 error when panic occurs inside the request call -func HandlePanic(next http.Handler) http.Handler { +func handlePanic(next http.Handler) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, req *http.Request) { defer func() { if err := recover(); err != nil { diff --git a/pkg/api/server.go b/pkg/api/server.go index b68d2e171..7852dac89 100644 --- a/pkg/api/server.go +++ b/pkg/api/server.go @@ -4,81 +4,40 @@ import ( "context" "net" "net/http" - "net/url" - "os" - "strings" - "sync" - "time" "github.com/go-chi/chi" - "github.com/pkg/errors" - "github.com/spf13/cobra" + "github.com/go-chi/chi/middleware" "github.com/titpetric/factory/resputil" "go.uber.org/zap" + "github.com/cortezaproject/corteza-server/pkg/app/options" "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli/options" "github.com/cortezaproject/corteza-server/pkg/version" ) type ( - Server struct { - name string - - log *zap.Logger - - httpOpt *options.HTTPOpt - monitorOpt *options.MonitorOpt - - endpoints []func(r chi.Router) + server struct { + log *zap.Logger + httpOpt options.HTTPServerOpt + waitForOpt options.WaitForOpt + endpoints []func(r chi.Router) } ) -var ( - Monolith = false - BaseURL = "/" -) - -func NewServer(log *zap.Logger) *Server { - return &Server{ - endpoints: make([]func(r chi.Router), 0), - log: log.Named("http"), +func New(log *zap.Logger, httpOpt options.HTTPServerOpt, waitForOpt options.WaitForOpt) *server { + return &server{ + endpoints: make([]func(r chi.Router), 0), + log: log.Named("http"), + httpOpt: httpOpt, + waitForOpt: waitForOpt, } } -func (s *Server) Command(ctx context.Context, cmdName, prefix string, preRun func(context.Context) error) (cmd *cobra.Command) { - s.httpOpt = options.HTTP(prefix) - s.monitorOpt = options.Monitor(prefix) - - cmd = &cobra.Command{ - Use: cmdName, - Short: "Start HTTP Server with REST API", - - // Connect all the wires, prepare services, run watchers, bind endpoints - PreRunE: func(cmd *cobra.Command, args []string) error { - s.waitFor(ctx, options.WaitFor(prefix)) - - if s.monitorOpt.Interval > 0 { - go NewMonitor(int(s.monitorOpt.Interval / time.Second)) - } - - return preRun(ctx) - }, - - // Run the server - Run: func(cmd *cobra.Command, args []string) { - s.Serve(ctx) - }, - } - - return -} - -func (s *Server) MountRoutes(mm ...func(chi.Router)) { +func (s *server) MountRoutes(mm ...func(chi.Router)) { s.endpoints = append(s.endpoints, mm...) } -func (s Server) Serve(ctx context.Context) { +func (s server) Serve(ctx context.Context) { s.log.Info("Starting HTTP server with REST API", zap.String("address", s.httpOpt.Addr)) // configure resputil options @@ -98,7 +57,7 @@ func (s Server) Serve(ctx context.Context) { router := chi.NewRouter() // Base middleware, CORS, RealIP, RequestID, context-logger - router.Use(Base(s.log)...) + router.Use(BaseMiddleware(s.log)...) // Logging request if enabled if s.httpOpt.LogRequest { @@ -110,17 +69,17 @@ func (s Server) Serve(ctx context.Context) { router.Use(LogResponse) } - // Handle panic (sets 500 Server error headers) - router.Use(HandlePanic) + // Handle panic (sets 500 server error headers) + router.Use(handlePanic) // Reports error to Sentry if enabled if s.httpOpt.EnablePanicReporting { - router.Use(Sentry()) + router.Use(sentryMiddleware()) } // Metrics tracking middleware if s.httpOpt.EnableMetrics { - router.Use(Middleware(s.httpOpt.MetricsServiceLabel)) + router.Use(metricsMiddleware(s.httpOpt.MetricsServiceLabel)) } router.Group(func(r chi.Router) { @@ -135,11 +94,15 @@ func (s Server) Serve(ctx context.Context) { }) if s.httpOpt.EnableMetrics { - Mount(router, s.httpOpt.MetricsUsername, s.httpOpt.MetricsPassword) + metricsMount(router, s.httpOpt.MetricsUsername, s.httpOpt.MetricsPassword) } if s.httpOpt.EnableDebugRoute { - Debug(router) + s.log.Debug("profiler: /__profiler", zap.Error(err)) + router.Mount("/__profiler", middleware.Profiler()) + + s.log.Debug("list of routes: /__routes", zap.Error(err)) + router.Get("/__routes", debugRoutes(router)) } if s.httpOpt.EnableVersionRoute { @@ -158,171 +121,5 @@ func (s Server) Serve(ctx context.Context) { } } - s.log.Info("HTTP server stopped", zap.Error(err)) - - return -} - -// waitFor sets up a simple status page, delays execution and probes services -func (s Server) waitFor(ctx context.Context, opt *options.WaitForOpt) { - var ( - services = opt.GetServices() - ) - - if len(services) == 0 && opt.Delay == 0 { - // Nothing to do here.. - return - } - - var ( - log = s.log.Named("wait-for") - depChan = make(chan struct{}) - wg sync.WaitGroup - serviceAddr string - serviceURL *url.URL - err error - ) - - // Setup a simple HTTP server that will inform the impatent users - listener, err := net.Listen("tcp", s.httpOpt.Addr) - if err != nil { - s.log.Error("Can not start server", zap.Error(err)) - os.Exit(1) - } - defer listener.Close() - go func() { - router := chi.NewRouter() - router.Get("/*", func(w http.ResponseWriter, r *http.Request) { - w.WriteHeader(http.StatusPreconditionFailed) - w.Write([]byte("waiting for services...")) - }) - _ = http.Serve(listener, router) - }() - - if opt.Delay > 0 { - s.log.Info("delaying", zap.Duration("delay", opt.Delay)) - - // First delay execution - select { - case <-ctx.Done(): - log.Debug("canceled") - return - case <-time.After(opt.Delay): - // all good... - } - } - - if len(services) == 0 { - return - } - - log.Info("waiting for services", zap.Strings("services", services)) - // Probe services - wg.Add(len(services)) - go func() { - - for _, service := range services { - slog := log.With(zap.String("service", service)) - - go func(ctx context.Context, service string) { - defer wg.Done() - - if serviceAddr, serviceURL, err = s.resolveService(service); err != nil { - log.Error("could not resolve service", zap.Error(err)) - } - - for { - ctx, cancelFn := context.WithTimeout(ctx, opt.ServicesProbeTimeout) - defer cancelFn() - - if serviceURL == nil { - if err = s.probeService(ctx, serviceAddr); err != nil { - slog.Warn("service probe failed", zap.Error(err)) - time.Sleep(opt.ServicesProbeInterval) - continue - } - } else { - if err = s.probeServiceURL(ctx, serviceURL); err != nil { - slog.Warn("service URL probe failed", zap.Error(err)) - time.Sleep(opt.ServicesProbeInterval) - continue - } - } - - slog.Debug("service ready") - return - } - }(ctx, service) - } - wg.Wait() - close(depChan) - }() - - select { - case <-ctx.Done(): - log.Debug("canceled") - return - case <-depChan: // services are ready - log.Debug("all services ready") - return - case <-time.After(opt.ServicesTimeout): - log.Debug("services not ready") - os.Exit(1) - } -} - -func (s Server) resolveService(service string) (addr string, u *url.URL, err error) { - addr = service - - if strings.Contains(addr, "://") { - // Is service an URL? - u, err = url.Parse(addr) - if err != nil { - return - } - - addr = u.Host - - if u.Port() == "" { - if u.Scheme == "https" { - addr += ":443" - } - } - } - - // Default to port 80 - if !strings.Contains(addr, ":") { - addr += ":80" - } - - return -} - -func (s Server) probeService(ctx context.Context, addr string) (err error) { - if err != nil { - return err - } - - dialer := net.Dialer{} - _, err = dialer.DialContext(ctx, "tcp", addr) - return -} - -func (s Server) probeServiceURL(ctx context.Context, u *url.URL) error { - req, err := http.NewRequest("GET", u.String(), nil) - if err != nil { - return errors.Wrap(err, "failed to assemble service request") - } - - rsp, err := http.DefaultClient.Do(req.WithContext(ctx)) - if err != nil { - return errors.Wrap(err, "service URL request failed") - } - - defer rsp.Body.Close() - if rsp.StatusCode == http.StatusOK { - return nil - } - - return errors.Errorf("service responded with unexpected status '%s'", rsp.Status) + s.log.Info("Server stopped", zap.Error(err)) } diff --git a/pkg/api/waitfor.go b/pkg/api/waitfor.go new file mode 100644 index 000000000..c868946f7 --- /dev/null +++ b/pkg/api/waitfor.go @@ -0,0 +1,181 @@ +package api + +import ( + "context" + "net" + "net/http" + "net/url" + "os" + "strings" + "sync" + "time" + + "github.com/go-chi/chi" + "github.com/pkg/errors" + "go.uber.org/zap" +) + +// WaitFor sets up a simple status page, delays execution and probes services +func (s server) WaitFor(ctx context.Context) { + var ( + opt = s.waitForOpt + services = opt.GetServices() + ) + + if len(services) == 0 && opt.Delay == 0 { + // Nothing to do here.. + return + } + + var ( + log = s.log.Named("wait-for") + depChan = make(chan struct{}) + wg sync.WaitGroup + serviceAddr string + serviceURL *url.URL + err error + ) + + // Setup a simple HTTP server that will inform the impatient users + listener, err := net.Listen("tcp", s.httpOpt.Addr) + if err != nil { + s.log.Error("Can not start server", zap.Error(err)) + os.Exit(1) + } + defer listener.Close() + go func() { + router := chi.NewRouter() + router.Get("/*", func(w http.ResponseWriter, r *http.Request) { + w.WriteHeader(http.StatusPreconditionFailed) + w.Write([]byte("waiting for services...")) + }) + _ = http.Serve(listener, router) + }() + + if opt.Delay > 0 { + s.log.Info("delaying", zap.Duration("delay", opt.Delay)) + + // First delay execution + select { + case <-ctx.Done(): + log.Debug("canceled") + return + case <-time.After(opt.Delay): + // all good... + } + } + + if len(services) == 0 { + return + } + + log.Info("waiting for services", zap.Strings("services", services)) + // Probe services + wg.Add(len(services)) + go func() { + + for _, service := range services { + slog := log.With(zap.String("service", service)) + + go func(ctx context.Context, service string) { + defer wg.Done() + + if serviceAddr, serviceURL, err = s.resolveService(service); err != nil { + log.Error("could not resolve service", zap.Error(err)) + } + + for { + ctx, cancelFn := context.WithTimeout(ctx, opt.ServicesProbeTimeout) + defer cancelFn() + + if serviceURL == nil { + if err = s.probeService(ctx, serviceAddr); err != nil { + slog.Warn("service probe failed", zap.Error(err)) + time.Sleep(opt.ServicesProbeInterval) + continue + } + } else { + if err = s.probeServiceURL(ctx, serviceURL); err != nil { + slog.Warn("service URL probe failed", zap.Error(err)) + time.Sleep(opt.ServicesProbeInterval) + continue + } + } + + slog.Debug("service ready") + return + } + }(ctx, service) + } + wg.Wait() + close(depChan) + }() + + select { + case <-ctx.Done(): + log.Debug("canceled") + return + case <-depChan: // services are ready + log.Debug("all services ready") + return + case <-time.After(opt.ServicesTimeout): + log.Debug("services not ready") + os.Exit(1) + } +} + +func (s server) resolveService(service string) (addr string, u *url.URL, err error) { + addr = service + + if strings.Contains(addr, "://") { + // Is service an URL? + u, err = url.Parse(addr) + if err != nil { + return + } + + addr = u.Host + + if u.Port() == "" { + if u.Scheme == "https" { + addr += ":443" + } + } + } + + // Default to port 80 + if !strings.Contains(addr, ":") { + addr += ":80" + } + + return +} + +func (s server) probeService(ctx context.Context, addr string) (err error) { + if err != nil { + return err + } + + dialer := net.Dialer{} + _, err = dialer.DialContext(ctx, "tcp", addr) + return +} + +func (s server) probeServiceURL(ctx context.Context, u *url.URL) error { + req, err := http.NewRequest("GET", u.String(), nil) + if err != nil { + return errors.Wrap(err, "failed to assemble service request") + } + + rsp, err := http.DefaultClient.Do(req.WithContext(ctx)) + if err != nil { + return errors.Wrap(err, "service URL request failed") + } + + defer rsp.Body.Close() + if rsp.StatusCode == http.StatusOK { + return nil + } + + return errors.Errorf("service responded with unexpected status '%s'", rsp.Status) +} diff --git a/pkg/app/helpers.go b/pkg/app/helpers.go new file mode 100644 index 000000000..2aaea2535 --- /dev/null +++ b/pkg/app/helpers.go @@ -0,0 +1,77 @@ +package app + +import ( + "context" + + "go.uber.org/zap" +) + +// RunSetup calls Setup hooks on all runnable parts +// +// It stops on first error +func RunSetup(log *zap.Logger, opts *Options, pp ...Runnable) (err error) { + for _, app := range pp { + err = app.Setup(log, opts) + if err != nil { + return + } + } + + return +} + +// RunInitialize calls Initialize hooks on all runnable parts +// +// It stops on first error +func RunInitialize(ctx context.Context, pp ...Runnable) (err error) { + for _, app := range pp { + err = app.Initialize(ctx) + if err != nil { + return + } + } + + return +} + +// RunUpgrade calls Upgrade hooks on all runnable parts +// +// It stops on first error +func RunUpgrade(ctx context.Context, pp ...Runnable) (err error) { + for _, app := range pp { + err = app.Upgrade(ctx) + if err != nil { + return + } + } + + return +} + +// RunActivate calls Activate hooks on all runnable parts +// +// It stops on first error +func RunActivate(ctx context.Context, pp ...Runnable) (err error) { + for _, app := range pp { + err = app.Activate(ctx) + if err != nil { + return + } + } + + return +} + +// RunProvision calls Provision hooks on all runnable parts +// +// It stops on first error +func RunProvision(ctx context.Context, pp ...Runnable) (err error) { + for _, app := range pp { + err = app.Provision(ctx) + if err != nil { + return + } + } + + return +} diff --git a/pkg/app/options.go b/pkg/app/options.go new file mode 100644 index 000000000..b80e7a50a --- /dev/null +++ b/pkg/app/options.go @@ -0,0 +1,48 @@ +package app + +import ( + "github.com/cortezaproject/corteza-server/pkg/app/options" +) + +type ( + Options struct { + SMTP options.SMTPOpt + JWT options.JWTOpt + HTTPClient options.HTTPClientOpt + DB options.DBOpt + Upgrade options.UpgradeOpt + Provision options.ProvisionOpt + Sentry options.SentryOpt + Storage options.StorageOpt + Corredor options.CorredorOpt + Monitor options.MonitorOpt + WaitFor options.WaitForOpt + HTTPServer options.HTTPServerOpt + GRPCServer options.GRPCServerOpt + Websocket options.WebsocketOpt + } +) + +func NewOptions(prefix ...string) *Options { + var p = "" + if len(prefix) > 0 { + p = prefix[0] + } + + return &Options{ + SMTP: *options.SMTP(p), + JWT: *options.JWT(p), + HTTPClient: *options.HttpClient(p), + DB: *options.DB(p), + Upgrade: *options.Upgrade(p), + Provision: *options.Provision(p), + Sentry: *options.Sentry(p), + Storage: *options.Storage(p), + Corredor: *options.Corredor(p), + Monitor: *options.Monitor(p), + WaitFor: *options.WaitFor(p), + HTTPServer: *options.HTTP(p), + GRPCServer: *options.GRPCServer(p), + Websocket: *options.Websocket(p), + } +} diff --git a/pkg/cli/options/corredor.go b/pkg/app/options/corredor.go similarity index 68% rename from pkg/cli/options/corredor.go rename to pkg/app/options/corredor.go index ea352e377..ac6aaffc2 100644 --- a/pkg/cli/options/corredor.go +++ b/pkg/app/options/corredor.go @@ -15,11 +15,6 @@ type ( Log bool `env:"CORREDOR_LOG_ENABLED"` MaxBackoffDelay time.Duration `env:"CORREDOR_MAX_BACKOFF_DELAY"` - - // @todo do autodiscovery & prefill these with values we know - ApiBaseURLSystem string `env:"CORREDOR_API_BASE_URL_SYSTEM"` - ApiBaseURLMessaging string `env:"CORREDOR_API_BASE_URL_MESSAGING"` - ApiBaseURLCompose string `env:"CORREDOR_API_BASE_URL_COMPOSE"` } ) diff --git a/pkg/cli/options/db.go b/pkg/app/options/db.go similarity index 100% rename from pkg/cli/options/db.go rename to pkg/app/options/db.go diff --git a/pkg/cli/options/grpc_server.go b/pkg/app/options/grpc_server.go similarity index 100% rename from pkg/cli/options/grpc_server.go rename to pkg/app/options/grpc_server.go diff --git a/pkg/cli/options/helpers.go b/pkg/app/options/helpers.go similarity index 100% rename from pkg/cli/options/helpers.go rename to pkg/app/options/helpers.go diff --git a/pkg/cli/options/http.go b/pkg/app/options/http.go similarity index 93% rename from pkg/cli/options/http.go rename to pkg/app/options/http.go index d781746fc..7d5c65694 100644 --- a/pkg/cli/options/http.go +++ b/pkg/app/options/http.go @@ -5,7 +5,7 @@ import ( ) type ( - HTTPOpt struct { + HTTPServerOpt struct { Addr string `env:"HTTP_ADDR"` LogRequest bool `env:"HTTP_LOG_REQUEST"` LogResponse bool `env:"HTTP_LOG_RESPONSE"` @@ -23,8 +23,8 @@ type ( } ) -func HTTP(pfix string) (o *HTTPOpt) { - o = &HTTPOpt{ +func HTTP(pfix string) (o *HTTPServerOpt) { + o = &HTTPServerOpt{ Addr: ":80", LogRequest: false, LogResponse: false, diff --git a/pkg/cli/options/http_client.go b/pkg/app/options/http_client.go similarity index 74% rename from pkg/cli/options/http_client.go rename to pkg/app/options/http_client.go index ee8493675..8b17edec9 100644 --- a/pkg/cli/options/http_client.go +++ b/pkg/app/options/http_client.go @@ -5,14 +5,14 @@ import ( ) type ( - HttpClientOpt struct { + HTTPClientOpt struct { ClientTSLInsecure bool `env:"HTTP_CLIENT_TSL_INSECURE"` HttpClientTimeout time.Duration `env:"HTTP_CLIENT_TIMEOUT"` } ) -func HttpClient(pfix string) (o *HttpClientOpt) { - o = &HttpClientOpt{ +func HttpClient(pfix string) (o *HTTPClientOpt) { + o = &HTTPClientOpt{ ClientTSLInsecure: false, HttpClientTimeout: 30 * time.Second, } diff --git a/pkg/cli/options/jwt.go b/pkg/app/options/jwt.go similarity index 100% rename from pkg/cli/options/jwt.go rename to pkg/app/options/jwt.go diff --git a/pkg/cli/options/monitor.go b/pkg/app/options/monitor.go similarity index 100% rename from pkg/cli/options/monitor.go rename to pkg/app/options/monitor.go diff --git a/pkg/app/options/provision.go b/pkg/app/options/provision.go new file mode 100644 index 000000000..199d010c7 --- /dev/null +++ b/pkg/app/options/provision.go @@ -0,0 +1,17 @@ +package options + +type ( + ProvisionOpt struct { + Always bool `env:"PROVISION_ALWAYS"` + } +) + +func Provision(pfix string) (o *ProvisionOpt) { + o = &ProvisionOpt{ + Always: true, + } + + fill(o, pfix) + + return +} diff --git a/pkg/cli/options/pubsub.go b/pkg/app/options/pubsub.go similarity index 100% rename from pkg/cli/options/pubsub.go rename to pkg/app/options/pubsub.go diff --git a/pkg/cli/options/sentry.go b/pkg/app/options/sentry.go similarity index 100% rename from pkg/cli/options/sentry.go rename to pkg/app/options/sentry.go diff --git a/pkg/cli/options/smtp.go b/pkg/app/options/smtp.go similarity index 100% rename from pkg/cli/options/smtp.go rename to pkg/app/options/smtp.go diff --git a/pkg/cli/options/storage.go b/pkg/app/options/storage.go similarity index 100% rename from pkg/cli/options/storage.go rename to pkg/app/options/storage.go diff --git a/pkg/app/options/upgrade.go b/pkg/app/options/upgrade.go new file mode 100644 index 000000000..6cac9acd0 --- /dev/null +++ b/pkg/app/options/upgrade.go @@ -0,0 +1,17 @@ +package options + +type ( + UpgradeOpt struct { + Always bool `env:"UPGRADE_ALWAYS"` + } +) + +func Upgrade(pfix string) (o *UpgradeOpt) { + o = &UpgradeOpt{ + Always: true, + } + + fill(o, pfix) + + return +} diff --git a/pkg/cli/options/wait_for.go b/pkg/app/options/wait_for.go similarity index 100% rename from pkg/cli/options/wait_for.go rename to pkg/app/options/wait_for.go diff --git a/pkg/cli/options/websocket.go b/pkg/app/options/websocket.go similarity index 100% rename from pkg/cli/options/websocket.go rename to pkg/app/options/websocket.go diff --git a/pkg/app/runner.go b/pkg/app/runner.go new file mode 100644 index 000000000..03c415ce6 --- /dev/null +++ b/pkg/app/runner.go @@ -0,0 +1,327 @@ +package app + +import ( + "context" + "sync" + + "github.com/go-chi/chi" + "github.com/spf13/cobra" + "go.uber.org/zap" + "google.golang.org/grpc" + + "github.com/cortezaproject/corteza-server/pkg/api" + "github.com/cortezaproject/corteza-server/pkg/cli" + grpcWrap "github.com/cortezaproject/corteza-server/pkg/grpc" +) + +type ( + runner struct { + execStep int + + log *zap.Logger + opt *Options + + parts []Runnable + + // CLI Commands + rootCmd *cobra.Command + + // Servers + httpApiServer httpApiServer + grpcServer grpcServer + } + + httpApiServer interface { + MountRoutes(mm ...func(chi.Router)) + Serve(context.Context) + } + + grpcServer interface { + RegisterServices(func(*grpc.Server)) + Serve(ctx context.Context) + } + + Runnable interface { + // Setup + // + // Step 0 + Setup(log *zap.Logger, opts *Options) error + + // Initialize initializes all (passive) services + // + // Step 1 + Initialize(ctx context.Context) error + + // Upgrade runs database migrations + // dep: Initialize + // + // Step 2 + Upgrade(ctx context.Context) error + + // Activate wakes-up all (active) services, watchers... + // dep: Upgrade + // + // Step 3 + Activate(ctx context.Context) error + + // Provision fills database with data + // dep: Activate + // + // Ran before each cli command (and before REST API) + // + // Step 4 + Provision(ctx context.Context) error + } + + cliCommandRegistrator interface { + // RegisterCliCommands registers all command line interface commands + RegisterCliCommands(cmd *cobra.Command) + } + + apiRouteMounter interface { + // MountApiRoutes mounts all routes + MountApiRoutes(router chi.Router) + } + + grpcRegistrator interface { + // RegisterGrpcServices registers all gRPC services that are needed + RegisterGrpcServices(server *grpc.Server) + } +) + +const ( + execStepReady = iota + execStepSetupDone + execStepInitializeDone + execStepUpgradeDone + execStepActivateDone + execStepProvisionDone +) + +func New(parts ...Runnable) *runner { + return &runner{ + execStep: execStepReady, + parts: parts, + } +} + +// Setup runs setup on all parts +func (r *runner) Setup(log *zap.Logger, opt *Options) (err error) { + r.log = log + r.opt = opt + return r.setup() +} + +func (r *runner) setup() (err error) { + if r.execStep < execStepSetupDone { + if err = RunSetup(r.log, r.opt, r.parts...); err != nil { + return + } + + r.execStep = execStepSetupDone + + } + + return nil +} + +// Initialize setups and initializes all parts +func (r *runner) Initialize(ctx context.Context) (err error) { + if r.execStep < execStepInitializeDone { + if err = r.setup(); err != nil { + return + } + + if err = RunInitialize(ctx, r.parts...); err != nil { + return + } + + r.execStep = execStepInitializeDone + } + + return nil + +} + +// Upgrade initializes & upgrades all parts +func (r *runner) Upgrade(ctx context.Context) (err error) { + if r.execStep < execStepUpgradeDone { + if err = r.Initialize(ctx); err != nil { + return + } + + if r.opt.Upgrade.Always { + if err = RunUpgrade(ctx, r.parts...); err != nil { + return + } + } + + r.execStep = execStepUpgradeDone + } + + return nil + +} + +// Activate upgrades and activates all parts +func (r *runner) Activate(ctx context.Context) (err error) { + if r.execStep < execStepActivateDone { + if err = r.Upgrade(ctx); err != nil { + return + } + + if err = RunActivate(ctx, r.parts...); err != nil { + return + } + + r.execStep = execStepActivateDone + } + + return nil +} + +// Provision activates and provisions all parts +func (r *runner) Provision(ctx context.Context) (err error) { + if r.execStep < execStepProvisionDone { + if err = r.Activate(ctx); err != nil { + return + } + + if r.opt.Provision.Always { + if err = RunProvision(ctx, r.parts...); err != nil { + return + } + } + + r.execStep = execStepProvisionDone + } + + return nil +} + +func (r *runner) setupHttpApi() { + // @todo refactor wait-for out of HTTP API server. + r.httpApiServer = api.New(r.log, r.opt.HTTPServer, r.opt.WaitFor) + + // Mount all HTTP API endpoints + for _, part := range r.parts { + if reg, is := part.(apiRouteMounter); is { + r.httpApiServer.MountRoutes(reg.MountApiRoutes) + } + } +} + +func (r *runner) setupGRPCServices() { + + r.grpcServer = grpcWrap.New(r.log, r.opt.GRPCServer) + + var hasServices bool + + // Register GRPC services + for _, part := range r.parts { + if reg, is := part.(grpcRegistrator); is { + r.grpcServer.RegisterServices(reg.RegisterGrpcServices) + hasServices = true + } + } + + if !hasServices { + r.grpcServer = nil + } +} + +// serve starts all servers (HTTP API, GRPC) +func (r *runner) serve(ctx context.Context) (err error) { + if err = r.Provision(ctx); err != nil { + return + } + r.setupHttpApi() + r.setupGRPCServices() + + wg := &sync.WaitGroup{} + + { + wg.Add(1) + go func(ctx context.Context) { + r.httpApiServer.Serve(ctx) + wg.Done() + }(ctx) + } + + if r.grpcServer != nil { + wg.Add(1) + go func(ctx context.Context) { + r.grpcServer.Serve(ctx) + wg.Done() + }(ctx) + + } + + // Wait for all servers to be done + wg.Wait() + + return nil +} + +// Run will orchestrate all hooks, register commands, server and execute root command +func (r *runner) Run(ctx context.Context) error { + r.rootCmd = cli.RootCommand(func() error { + // Run initialization (and all steps before that) + return r.Initialize(ctx) + }) + + serveCmd := cli.ServeCommand(func() error { + return r.serve(ctx) + }) + + upgradeCmd := cli.UpgradeCommand(func() (err error) { + if err = r.Initialize(ctx); err != nil { + return + } + + if err = RunUpgrade(ctx, r.parts...); err != nil { + return + } + + return + }) + + provisionCmd := cli.ProvisionCommand(func() (err error) { + if err = r.Activate(ctx); err != nil { + return + } + + if err = RunProvision(ctx, r.parts...); err != nil { + return + } + + return + }) + + r.rootCmd.AddCommand( + serveCmd, + upgradeCmd, + provisionCmd, + ) + + // Register CLI commands from all parts (when compatible) + for _, part := range r.parts { + if reg, is := part.(cliCommandRegistrator); is { + reg.RegisterCliCommands(r.rootCmd) + } + } + + return r.rootCmd.Execute() +} + +// Run is an application runner +// +// It accepts set of Runnables and calls hooks on each +func Run(log *zap.Logger, opt *Options, parts ...Runnable) { + r := New(parts...) + r.Setup(log, opt) + + ctx := cli.Context() + + cli.HandleError(r.Run(ctx)) +} diff --git a/pkg/cli/commands.go b/pkg/cli/commands.go new file mode 100644 index 000000000..eddecc3f8 --- /dev/null +++ b/pkg/cli/commands.go @@ -0,0 +1,82 @@ +package cli + +import ( + "github.com/spf13/cobra" +) + +var ( + rootCommand = cobra.Command{ + Use: "corteza-server", + Aliases: []string{"corteza", "server"}, + TraverseChildren: true, + } + + serveApiCommand = cobra.Command{ + Use: "serve-api", + Aliases: []string{"serve"}, + + Short: "Start HTTP server with REST API", + } + + upgradeCommand = cobra.Command{ + Use: "upgrade", + Short: "Upgrade tasks", + } + + provisionCommand = cobra.Command{ + Use: "provision", + Short: "Provision tasks", + } +) + +// RootCommand creates root command with a simple persistent-pre-run callback +// +// Callback is called when not executed with help subcommand +func RootCommand(ppRunEfn func() error) *cobra.Command { + // Make a copy + var cmd = rootCommand + + cmd.PersistentPreRunE = func(cmd *cobra.Command, args []string) (err error) { + if cmd.Name() == "help" { + // Do not run hooks on parts on help command + return nil + } + + return ppRunEfn() + } + + return &cmd +} + +func ServeCommand(runEfn func() error) *cobra.Command { + // Make a copy + var cmd = serveApiCommand + + cmd.RunE = func(cmd *cobra.Command, args []string) (err error) { + return runEfn() + } + + return &cmd +} + +func UpgradeCommand(dbfn func() error) *cobra.Command { + // Make a copies + var dbCmd = upgradeCommand + + dbCmd.RunE = func(cmd *cobra.Command, args []string) (err error) { + return dbfn() + } + + return &dbCmd +} + +func ProvisionCommand(cffn func() error) *cobra.Command { + // Make a copies + var cfCmd = provisionCommand + + cfCmd.RunE = func(cmd *cobra.Command, args []string) (err error) { + return cffn() + } + + return &cfCmd +} diff --git a/pkg/cli/helpers.go b/pkg/cli/helpers.go index bbf76f061..f3240b715 100644 --- a/pkg/cli/helpers.go +++ b/pkg/cli/helpers.go @@ -3,23 +3,8 @@ package cli import ( "fmt" "os" - "time" - - "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli/options" - "github.com/cortezaproject/corteza-server/pkg/http" - "github.com/cortezaproject/corteza-server/pkg/mail" ) -func InitGeneralServices(smtpOpt *options.SMTPOpt, jwtOpt *options.JWTOpt, httpClientOpt *options.HttpClientOpt) { - auth.SetupDefault(jwtOpt.Secret, int(jwtOpt.Expiry/time.Minute)) - mail.SetupDialer(smtpOpt.Host, smtpOpt.Port, smtpOpt.User, smtpOpt.Pass, smtpOpt.From) - http.SetupDefaults( - httpClientOpt.HttpClientTimeout, - httpClientOpt.ClientTSLInsecure, - ) -} - func HandleError(err error) { if err == nil { return diff --git a/pkg/cli/options/provision.go b/pkg/cli/options/provision.go deleted file mode 100644 index 02fb024e3..000000000 --- a/pkg/cli/options/provision.go +++ /dev/null @@ -1,19 +0,0 @@ -package options - -type ( - ProvisionOpt struct { - MigrateDatabase bool `env:"PROVISION_MIGRATE_DATABASE"` - Configuration bool `env:"PROVISION_CONFIGURATION"` - } -) - -func Provision(pfix string) (o *ProvisionOpt) { - o = &ProvisionOpt{ - MigrateDatabase: true, - Configuration: true, - } - - fill(o, pfix) - - return -} diff --git a/pkg/cli/runner.go b/pkg/cli/runner.go deleted file mode 100644 index 07d76d48c..000000000 --- a/pkg/cli/runner.go +++ /dev/null @@ -1,305 +0,0 @@ -package cli - -import ( - "context" - - "github.com/go-chi/chi" - "github.com/pkg/errors" - "github.com/spf13/cobra" - "go.uber.org/zap" - - "github.com/cortezaproject/corteza-server/pkg/api" - "github.com/cortezaproject/corteza-server/pkg/cli/options" - "github.com/cortezaproject/corteza-server/pkg/db" - "github.com/cortezaproject/corteza-server/pkg/logger" - "github.com/cortezaproject/corteza-server/pkg/sentry" -) - -type ( - Runner func(ctx context.Context, cmd *cobra.Command, c *Config) error - Runners []Runner - - CommandMaker func(ctx context.Context, c *Config) *cobra.Command - CommandMakers []CommandMaker - - FlagBinder func(cmd *cobra.Command, c *Config) - FlagBinders []FlagBinder - - Mounter func(r chi.Router) - Mounters []Mounter - - Config struct { - init bool - - // Service name (messaging, system...) - // See comments on other fields for how it is used. - ServiceName string - - // Prefix for ENV variables - EnvPrefix string - - // Logger name for internal services, defaults to ServiceName - LoggerName string - Log *zap.Logger - - // General options - SmtpOpt *options.SMTPOpt - JwtOpt *options.JWTOpt - HttpClientOpt *options.HttpClientOpt - DbOpt *options.DBOpt - ProvisionOpt *options.ProvisionOpt - SentryOpt *options.SentryOpt - StorageOpt *options.StorageOpt - ScriptRunner *options.CorredorOpt - - // Services will be calling each other so we need - // to keep the config opts spearated - GRPCServerSystem *options.GRPCServerOpt - // GRPCServerMessaging *options.GRPCServerOpt - // GRPCServerCompose *options.GRPCServerOpt - - // DB Connection name, defaults to ServiceName - DatabaseName string - - // Root command name, , defaults to "corteza-server-" - RootCommandName string - - // Database setup/connection procedure - // Runner autobinds default runner that tries to connect using DbOpt.DSN - RootCommandDBSetup Runners - - // All that needs to be initialized before any sub-comman is executed - RootCommandPreRun Runners - - // ****************************************************************** - - // API Server instance - ApiServer *api.Server - - // API Server command name - ApiServerCommandName string - - // Code that needs to be executed before HTTP server is started - ApiServerPreRun Runners - - // Routers that we mount on HTTP server - ApiServerRoutes Mounters - - // Sets-up all available subcommands. - AdtSubCommands CommandMakers - - // Database migration code - // This is used for "provision migrate-database" command and after db connection is - // established (if --provision-migrate-database is enabled - ProvisionMigrateDatabase Runners - - // Access control initial setup - // Reapplies default access control rules for roles "everyone" [1] and "admin" [2] - ProvisionConfig Runners - - // ****************************************************************** - - // This callback behaves a bit differently and should be called manually - // from wherever we need service initialized - InitServices func(ctx context.Context, c *Config) - } -) - -func init() { - // Have logger ready in case we need to log anything - // before it gets properly initialized through InitGeneralServices - logger.Init() -} - -func (rr Runners) Run(ctx context.Context, cmd *cobra.Command, c *Config) (err error) { - for i := range rr { - err = rr[i](ctx, cmd, c) - if err != nil { - return - } - } - - return -} - -func (rr Mounters) MountRoutes(r chi.Router) { - for i := range rr { - rr[i](r) - } -} - -func (bb FlagBinders) Bind(cmd *cobra.Command, c *Config) { - for i := range bb { - bb[i](cmd, c) - } -} - -func (mm CommandMakers) Make(ctx context.Context, c *Config) []*cobra.Command { - var ( - valid = make([]*cobra.Command, 0) - cmd *cobra.Command - ) - - for i := range mm { - if cmd = mm[i](ctx, c); cmd != nil { - valid = append(valid, cmd) - } - } - - return valid -} - -func CombineFlagBinders(rr ...FlagBinders) (out FlagBinders) { - for i := range rr { - out = append(out, rr[i]...) - } - - return out -} - -func (c *Config) Init() { - if c.init { - return - } - - if c.Log == nil { - c.Log = logger.Default() - } - - if c.LoggerName == "" { - c.LoggerName = c.ServiceName - } - - if c.EnvPrefix == "" { - c.EnvPrefix = c.ServiceName - } - - c.Log = c.Log.Named(c.LoggerName) - - if c.RootCommandName == "" { - c.RootCommandName = "corteza-server-" + c.ServiceName - } - - if c.ApiServerCommandName == "" { - c.ApiServerCommandName = "serve-api" - } - - if c.DatabaseName == "" { - c.DatabaseName = c.ServiceName - } - - c.SmtpOpt = options.SMTP(c.EnvPrefix) - c.JwtOpt = options.JWT(c.EnvPrefix) - c.HttpClientOpt = options.HttpClient(c.EnvPrefix) - c.DbOpt = options.DB(c.ServiceName) - c.ProvisionOpt = options.Provision(c.ServiceName) - c.SentryOpt = options.Sentry(c.EnvPrefix) - c.StorageOpt = options.Storage(c.EnvPrefix) - c.ScriptRunner = options.Corredor(c.EnvPrefix) - c.GRPCServerSystem = options.GRPCServer("system") - // c.GRPCServerCompose = options.GRPCServer("compose") - // c.GRPCServerMessagign = options.GRPCServer("messaging") - - if c.RootCommandDBSetup == nil { - c.RootCommandDBSetup = Runners{func(ctx context.Context, cmd *cobra.Command, c *Config) (err error) { - if c.DbOpt != nil { - _, err = db.TryToConnect(ctx, c.Log, c.DatabaseName, *c.DbOpt) - if err != nil { - return errors.Wrap(err, "could not connect to database") - } - } - - return - }} - } - - if c.ApiServer == nil { - c.ApiServer = api.NewServer(c.Log) - } - - for i := range c.ApiServerRoutes { - c.ApiServer.MountRoutes(c.ApiServerRoutes[i]) - } -} - -// MakeCLI creates command line interface -// -// It tries to construct "serve-api" and "provision" sub-commands -// if configured properly (see Config struct) -func (c *Config) MakeCLI(ctx context.Context) (cmd *cobra.Command) { - c.Init() - - cmd = &cobra.Command{ - Use: c.RootCommandName, - TraverseChildren: true, - PersistentPreRunE: func(cmd *cobra.Command, args []string) (err error) { - if err = sentry.Init(c.SentryOpt); err != nil { - c.Log.Error("could not initialize Sentry", zap.Error(err)) - } - - defer sentry.Recover() - - InitGeneralServices(c.SmtpOpt, c.JwtOpt, c.HttpClientOpt) - - err = c.RootCommandDBSetup.Run(ctx, cmd, c) - if err != nil { - c.Log.Error("Failed to connect to the database", zap.Error(err)) - return nil - } - - err = c.RootCommandPreRun.Run(ctx, cmd, c) - if err != nil { - c.Log.Error("Failed to run command pre-run scripts", zap.Error(err)) - return nil - } - - return nil - }, - } - - serveApiCmd := c.ApiServer.Command(ctx, c.ApiServerCommandName, c.EnvPrefix, func(ctx context.Context) (err error) { - defer sentry.Recover() - return c.ApiServerPreRun.Run(ctx, cmd, c) - }) - - cmd.AddCommand(serveApiCmd) - - if len(c.ProvisionMigrateDatabase) > 0 || len(c.ProvisionConfig) > 0 { - var ( - provisionCmd = &cobra.Command{ - Use: "provision", - Short: "Provision tasks", - } - ) - - // Add only commands with defined callbacks - if len(c.ProvisionMigrateDatabase) > 0 { - provisionCmd.AddCommand(&cobra.Command{ - Use: "configuration", - Short: "Create permissions & resources", - - RunE: func(cmd *cobra.Command, args []string) error { - return c.ProvisionConfig.Run(ctx, nil, c) - }, - }) - } - - // Add only commands with defined callbacks - if len(c.ProvisionConfig) > 0 { - provisionCmd.AddCommand(&cobra.Command{ - Use: "migrate-database", - Short: "Run database migration scripts", - - RunE: func(cmd *cobra.Command, args []string) error { - return c.ProvisionMigrateDatabase.Run(ctx, nil, c) - }, - }) - } - - cmd.AddCommand(provisionCmd) - } - - cmd.AddCommand(c.AdtSubCommands.Make(ctx, c)...) - - return -} diff --git a/pkg/db/connector.go b/pkg/db/connector.go index 9666e3690..f1245da5a 100644 --- a/pkg/db/connector.go +++ b/pkg/db/connector.go @@ -10,7 +10,7 @@ import ( "github.com/titpetric/factory/logger" "go.uber.org/zap" - "github.com/cortezaproject/corteza-server/pkg/cli/options" + "github.com/cortezaproject/corteza-server/pkg/app/options" "github.com/cortezaproject/corteza-server/pkg/sentry" ) @@ -18,7 +18,13 @@ var ( dsnMasker = regexp.MustCompile("(.)(?:.*)(.):(.)(?:.*)(.)@") ) -func TryToConnect(ctx context.Context, log *zap.Logger, name string, opt options.DBOpt) (db *factory.DB, err error) { +func TryToConnect(ctx context.Context, log *zap.Logger, opt options.DBOpt) (db *factory.DB, err error) { + if opt.DSN == "" { + err = errors.Errorf("invalid or empty DSN: %q", opt.DSN) + return + } + + name := "default" factory.Database.Add(name, opt.DSN) var ( @@ -74,7 +80,7 @@ func TryToConnect(ctx context.Context, log *zap.Logger, name string, opt options } } - log.Info("connected to the database", dsnField) + log.Debug("connected to the database", dsnField) // Connected break diff --git a/system/grpc/server.go b/pkg/grpc/server.go similarity index 53% rename from system/grpc/server.go rename to pkg/grpc/server.go index a86df4b04..5bb657300 100644 --- a/system/grpc/server.go +++ b/pkg/grpc/server.go @@ -2,36 +2,69 @@ package grpc import ( "context" + "net" + + "go.uber.org/zap" "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/metadata" "google.golang.org/grpc/status" + "github.com/cortezaproject/corteza-server/pkg/app/options" "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/system/proto" - "github.com/cortezaproject/corteza-server/system/service" ) -// @todo when we extend gRPC-server capabilities to compose & messaging -// this needs to be refactored and generalized -func NewServer() *grpc.Server { - s := grpc.NewServer( - grpc.UnaryInterceptor(authCheck(auth.DefaultJwtHandler)), - ) +type ( + server struct { + log *zap.Logger + opt options.GRPCServerOpt - proto.RegisterUsersServer(s, NewUserService( - service.DefaultUser, - service.DefaultAuth, - auth.DefaultJwtHandler, - service.DefaultAccessControl, - )) + s *grpc.Server + } +) - proto.RegisterRolesServer(s, NewRoleService( - service.DefaultRole, - )) +func New(log *zap.Logger, opt options.GRPCServerOpt) *server { + return &server{ + log: log.Named("grpc-server"), + opt: opt, - return s + s: grpc.NewServer( + grpc.UnaryInterceptor(authCheck(auth.DefaultJwtHandler)), + ), + } +} + +func (srv *server) Serve(ctx context.Context) { + ln, err := net.Listen(srv.opt.Network, srv.opt.Addr) + if err != nil { + srv.log.Error("could not start gRPC server", zap.Error(err)) + } + + go func() { + select { + case <-ctx.Done(): + srv.log.Debug("shutting down") + srv.s.GracefulStop() + _ = ln.Close() + } + }() + + srv.log.Info("Starting gRPC server", zap.String("address", srv.opt.Addr)) + err = srv.s.Serve(ln) + + if err == nil { + err = ctx.Err() + if err == context.Canceled { + err = nil + } + } + + srv.log.Info("Server stopped", zap.Error(err)) +} + +func (srv *server) RegisterServices(reg func(*grpc.Server)) { + reg(srv.s) } // Creates auth-checking interceptor function diff --git a/pkg/logger/logger.go b/pkg/logger/logger.go index 640d34746..11df45851 100644 --- a/pkg/logger/logger.go +++ b/pkg/logger/logger.go @@ -2,6 +2,7 @@ package logger import ( "os" + "time" // Make sure we read the ENV from .env _ "github.com/joho/godotenv/autoload" @@ -9,7 +10,7 @@ import ( "go.uber.org/zap" "go.uber.org/zap/zapcore" - "github.com/cortezaproject/corteza-server/pkg/cli/options" + "github.com/cortezaproject/corteza-server/pkg/app/options" ) var ( @@ -21,6 +22,14 @@ func MakeDebugLogger() *zap.Logger { conf := zap.NewDevelopmentConfig() conf.Level = DefaultLevel + // Print log level in colors + conf.EncoderConfig.EncodeLevel = zapcore.CapitalColorLevelEncoder + + // Shorten timestamp, we do not care about the date + conf.EncoderConfig.EncodeTime = func(t time.Time, enc zapcore.PrimitiveArrayEncoder) { + enc.AppendString(t.Format("15:04:05.000")) + } + logger, err := conf.Build() if err != nil { panic(err) diff --git a/pkg/api/monitor.go b/pkg/monitor/monitor.go similarity index 53% rename from pkg/api/monitor.go rename to pkg/monitor/monitor.go index e56429b1d..823b62f02 100644 --- a/pkg/api/monitor.go +++ b/pkg/monitor/monitor.go @@ -1,13 +1,14 @@ -package api +package monitor import ( + "context" "expvar" "runtime" "time" "go.uber.org/zap" - "github.com/cortezaproject/corteza-server/pkg/logger" + "github.com/cortezaproject/corteza-server/pkg/app/options" ) type Monitor struct { @@ -23,6 +24,25 @@ type Monitor struct { NumGoroutine int } +var ( + // Holds options for monitor + opt options.MonitorOpt + + log *zap.Logger +) + +func Setup(logger *zap.Logger, o options.MonitorOpt) { + log = logger.Named("monitor") + opt = o +} + +func Watcher(ctx context.Context) { + if opt.Interval > 0 { + go NewMonitor(int(opt.Interval / time.Second)) + log.Debug("watcher initialized") + } +} + func NewMonitor(duration int) { var ( m = Monitor{} @@ -54,18 +74,17 @@ func NewMonitor(duration int) { m.PauseTotalNs = rtm.PauseTotalNs m.NumGC = rtm.NumGC - logger.Default(). - With( - zap.Uint64("alloc", m.Alloc), - zap.Uint64("totalAlloc", m.TotalAlloc), - zap.Uint64("sys", m.Sys), - zap.Uint64("mallocs", m.Mallocs), - zap.Uint64("frees", m.Frees), - zap.Uint64("liveObjects", m.LiveObjects), - zap.Uint64("pauseTotalNs", m.PauseTotalNs), - zap.Uint32("numGC", m.NumGC), - zap.Int("numGoRoutines", m.NumGoroutine), - ). - Debug("monitor") + log.With( + zap.Uint64("alloc", m.Alloc), + zap.Uint64("totalAlloc", m.TotalAlloc), + zap.Uint64("sys", m.Sys), + zap.Uint64("mallocs", m.Mallocs), + zap.Uint64("frees", m.Frees), + zap.Uint64("liveObjects", m.LiveObjects), + zap.Uint64("pauseTotalNs", m.PauseTotalNs), + zap.Uint32("numGC", m.NumGC), + zap.Int("numGoRoutines", m.NumGoroutine), + ). + Info("tick") } } diff --git a/pkg/sentry/sentry.go b/pkg/sentry/sentry.go index 1c62e10b4..be9eb2012 100644 --- a/pkg/sentry/sentry.go +++ b/pkg/sentry/sentry.go @@ -3,11 +3,11 @@ package sentry import ( "github.com/getsentry/sentry-go" - "github.com/cortezaproject/corteza-server/pkg/cli/options" + "github.com/cortezaproject/corteza-server/pkg/app/options" "github.com/cortezaproject/corteza-server/pkg/logger" ) -func Init(sentryOpt *options.SentryOpt) error { +func Init(sentryOpt options.SentryOpt) error { if sentryOpt.DSN == "" { return nil } diff --git a/pkg/settings/service.go b/pkg/settings/service.go index 9d681a600..72fd441a9 100644 --- a/pkg/settings/service.go +++ b/pkg/settings/service.go @@ -56,7 +56,6 @@ func (svc service) log(ctx context.Context, fields ...zapcore.Field) *zap.Logger func (svc service) FindByPrefix(ctx context.Context, pp ...string) (ValueSet, error) { if !svc.accessControl.CanReadSettings(ctx) { - svc.log(ctx).Error("foo") return nil, ErrNoReadPermission } diff --git a/system/app.go b/system/app.go new file mode 100644 index 000000000..93e2dd997 --- /dev/null +++ b/system/app.go @@ -0,0 +1,140 @@ +package system + +import ( + "context" + + "github.com/go-chi/chi" + _ "github.com/joho/godotenv/autoload" + "github.com/spf13/cobra" + "github.com/titpetric/factory" + "go.uber.org/zap" + "google.golang.org/grpc" + + "github.com/cortezaproject/corteza-server/pkg/app" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/system/auth/external" + "github.com/cortezaproject/corteza-server/system/commands" + migrate "github.com/cortezaproject/corteza-server/system/db" + grpc2 "github.com/cortezaproject/corteza-server/system/grpc" + "github.com/cortezaproject/corteza-server/system/proto" + "github.com/cortezaproject/corteza-server/system/rest" + "github.com/cortezaproject/corteza-server/system/service" +) + +type ( + App struct { + Opts *app.Options + Log *zap.Logger + } +) + +const SERVICE = "system" + +func (app *App) Setup(log *zap.Logger, opts *app.Options) (err error) { + app.Log = log.Named(SERVICE) + app.Opts = opts + return +} + +func (app *App) Upgrade(ctx context.Context) (err error) { + db := factory.Database.MustGet(SERVICE, "default").With(ctx).Quiet() + err = migrate.Migrate(db, app.Log) + if err != nil { + return + } + + return +} + +// Initialized +func (app *App) Initialize(ctx context.Context) (err error) { + // Connects to all services it needs to + err = service.Initialize(ctx, app.Log, service.Config{ + Storage: app.Opts.Storage, + Corredor: app.Opts.Corredor, + }) + + if err != nil { + return + } + + // Initialize external authentication (from default settings) + external.Init() + return +} + +func (app *App) Activate(ctx context.Context) (err error) { + if err = service.Activate(ctx); err != nil { + return + } + + service.Watchers(ctx) + + // Wire in internal JWT maker for Corredor + corredor.Service().SetJwtMaker(corredor.InternalAuthTokenMaker()) + + return +} + +func (app *App) Provision(ctx context.Context) (err error) { + ctx = auth.SetSuperUserContext(ctx) + + if err = provisionConfig(ctx, app.Log); err != nil { + return + } + + // creates default applications that will appear in Unify/One + // @todo migrate this to provisioning/YAML + if err = makeDefaultApplications(ctx, app.Log); err != nil { + return + } + + // auto-discovery auth.* settings + if err = authSettingsAutoDiscovery(ctx, app.Log, service.DefaultSettings); err != nil { + return + } + + // external provider auto configuration + // creates: auth.external.providers.(google|linkedin|github|facebook).* + if err = authAddExternals(ctx); err != nil { + return + } + + // OIDC provider auto configuration + // creates: auth.external.providers.openid-connect.* + if err = oidcAutoDiscovery(ctx, app.Log); err != nil { + return + } + + return +} + +func (app *App) MountApiRoutes(r chi.Router) { + rest.MountRoutes(r) +} + +func (app *App) RegisterGrpcServices(server *grpc.Server) { + proto.RegisterUsersServer(server, grpc2.NewUserService( + service.DefaultUser, + service.DefaultAuth, + auth.DefaultJwtHandler, + service.DefaultAccessControl, + )) + + proto.RegisterRolesServer(server, grpc2.NewRoleService( + service.DefaultRole, + )) +} + +func (app *App) RegisterCliCommands(p *cobra.Command) { + p.AddCommand( + commands.Importer(), + commands.Exporter(), + commands.Settings(), + commands.Auth(), + commands.Users(), + commands.Roles(), + commands.Sink(), + ) +} diff --git a/system/app_test.go b/system/app_test.go new file mode 100644 index 000000000..51c2b5b3e --- /dev/null +++ b/system/app_test.go @@ -0,0 +1,11 @@ +package system + +import ( + "testing" + + "github.com/cortezaproject/corteza-server/pkg/app" +) + +func TestConfigure(t *testing.T) { + var _ app.Runnable = &App{} +} diff --git a/system/commands/auth.go b/system/commands/auth.go index c0846fd28..75639c978 100644 --- a/system/commands/auth.go +++ b/system/commands/auth.go @@ -1,7 +1,6 @@ package commands import ( - "context" "regexp" "strconv" @@ -17,7 +16,7 @@ import ( ) // Will perform OpenID connect auto-configuration -func Auth(ctx context.Context, c *cli.Config) *cobra.Command { +func Auth() *cobra.Command { var ( enableDiscoveredProvider bool skipValidationOnAutoDiscoveredProvider bool @@ -33,7 +32,9 @@ func Auth(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Auto discovers new OIDC client", Args: cobra.ExactArgs(2), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) + var ( + ctx = auth.SetSuperUserContext(cli.Context()) + ) _, err := external.RegisterOidcProvider( ctx, @@ -71,9 +72,9 @@ func Auth(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Generates new JWT for a user", Args: cobra.MinimumNArgs(1), Run: func(cmd *cobra.Command, args []string) { - var ( - db = factory.Database.MustGet("system") + ctx = auth.SetSuperUserContext(cli.Context()) + db = factory.Database.MustGet("system", "default") userRepo = repository.User(ctx, db) roleRepo = repository.Role(ctx, db) @@ -87,8 +88,6 @@ func Auth(ctx context.Context, c *cli.Config) *cobra.Command { userStr = args[0] ) - c.InitServices(ctx, c) - if user, err = userRepo.FindByEmail(userStr); repository.ErrUserNotFound.Eq(err) { if regexp.MustCompile(`/^\d+$/`).MatchString(userStr) { if ID, err = strconv.ParseUint(userStr, 10, 64); err == nil { @@ -116,9 +115,8 @@ func Auth(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Sends samples of all authentication notification to receipient", Args: cobra.ExactArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) err error ntf = service.DefaultAuthNotification.With(ctx) ) diff --git a/system/commands/exporter.go b/system/commands/exporter.go index 5576fdcbe..075491381 100644 --- a/system/commands/exporter.go +++ b/system/commands/exporter.go @@ -16,18 +16,16 @@ import ( sysTypes "github.com/cortezaproject/corteza-server/system/types" ) -func Exporter(ctx context.Context, c *cli.Config) *cobra.Command { +func Exporter() *cobra.Command { cmd := &cobra.Command{ Use: "export", Short: "Export", Long: `Export system resources`, Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) + sFlag = cmd.Flags().Lookup("settings").Changed pFlag = cmd.Flags().Lookup("permissions").Changed diff --git a/system/commands/importer.go b/system/commands/importer.go index 684f43999..b98240fbb 100644 --- a/system/commands/importer.go +++ b/system/commands/importer.go @@ -1,7 +1,6 @@ package commands import ( - "context" "io" "os" @@ -12,22 +11,18 @@ import ( "github.com/cortezaproject/corteza-server/system/importer" ) -func Importer(ctx context.Context, c *cli.Config) *cobra.Command { +func Importer() *cobra.Command { cmd := &cobra.Command{ Use: "import", Short: "Import", Run: func(cmd *cobra.Command, args []string) { - - c.InitServices(ctx, c) - var ( ff []io.Reader err error + ctx = auth.SetSuperUserContext(cli.Context()) ) - ctx = auth.SetSuperUserContext(ctx) - if len(args) > 0 { ff = make([]io.Reader, len(args)) for a, arg := range args { diff --git a/system/commands/roles.go b/system/commands/roles.go index 0b5a42566..251ca0cde 100644 --- a/system/commands/roles.go +++ b/system/commands/roles.go @@ -1,19 +1,19 @@ package commands import ( - "context" "strconv" "github.com/pkg/errors" "github.com/spf13/cobra" "github.com/titpetric/factory" + "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/pkg/cli" "github.com/cortezaproject/corteza-server/system/repository" "github.com/cortezaproject/corteza-server/system/types" ) -func Roles(ctx context.Context, c *cli.Config) *cobra.Command { +func Roles() *cobra.Command { cmd := &cobra.Command{ Use: "roles", Short: "Role management", @@ -26,7 +26,8 @@ func Roles(ctx context.Context, c *cli.Config) *cobra.Command { Run: func(cmd *cobra.Command, args []string) { // Create role and user repository. var ( - db = factory.Database.MustGet("system") + ctx = auth.SetSuperUserContext(cli.Context()) + db = factory.Database.MustGet("system", "default") roleStr, userStr = args[0], args[1] @@ -41,8 +42,6 @@ func Roles(ctx context.Context, c *cli.Config) *cobra.Command { err error ) - c.InitServices(ctx, c) - // Try to find role by name and by ID if rr, _, err = roleRepo.Find(types.RoleFilter{Query: roleStr}); err != nil { cli.HandleError(err) diff --git a/system/commands/settings.go b/system/commands/settings.go index fed821933..bfc117e8a 100644 --- a/system/commands/settings.go +++ b/system/commands/settings.go @@ -1,7 +1,6 @@ package commands import ( - "context" "encoding/json" "github.com/cortezaproject/corteza-server/pkg/auth" "os" @@ -9,12 +8,13 @@ import ( "github.com/spf13/cobra" + "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/pkg/cli" "github.com/cortezaproject/corteza-server/pkg/settings" "github.com/cortezaproject/corteza-server/system/service" ) -func Settings(ctx context.Context, c *cli.Config) *cobra.Command { +func Settings() *cobra.Command { var ( cmd = &cobra.Command{ Use: "settings", @@ -26,8 +26,9 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Use: "list", Short: "List all", Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) + var ( + ctx = auth.SetSuperUserContext(cli.Context()) + ) prefix := cmd.Flags().Lookup("prefix").Value.String() if kv, err := service.DefaultSettings.FindByPrefix(ctx, prefix); err != nil { @@ -55,8 +56,9 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Get value (raw JSON) for a specific key", Args: cobra.ExactArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) + var ( + ctx = auth.SetSuperUserContext(cli.Context()) + ) if v, err := service.DefaultSettings.Get(ctx, args[0], 0); err != nil { cli.HandleError(err) @@ -71,8 +73,9 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Set value (raw JSON) for a specific key", Args: cobra.ExactArgs(2), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) + var ( + ctx = auth.SetSuperUserContext(cli.Context()) + ) value := args[1] v := &settings.Value{ @@ -98,10 +101,8 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Import settings as JSON from stdin or file", Args: cobra.MaximumNArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) fh *os.File err error ) @@ -139,10 +140,8 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Import settings as JSON to stdout or file", Args: cobra.MaximumNArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( + ctx = auth.SetSuperUserContext(cli.Context()) fh *os.File err error ) @@ -173,10 +172,10 @@ func Settings(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Set value (raw JSON) for a specific key (or by prefix)", Args: cobra.MinimumNArgs(0), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - - var names = []string{} + var ( + ctx = auth.SetSuperUserContext(cli.Context()) + names = []string{} + ) if prefix := cmd.Flags().Lookup("prefix").Value.String(); len(prefix) > 0 { if vv, err := service.DefaultSettings.FindByPrefix(ctx); err != nil { diff --git a/system/commands/sink.go b/system/commands/sink.go index 379ae8ae5..2f54c4621 100644 --- a/system/commands/sink.go +++ b/system/commands/sink.go @@ -1,7 +1,6 @@ package commands import ( - "context" "net/url" "strings" "time" @@ -9,11 +8,10 @@ import ( "github.com/spf13/cobra" "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli" ) // Will perform OpenID connect auto-configuration -func Sink(ctx context.Context, c *cli.Config) *cobra.Command { +func Sink() *cobra.Command { var ( expires string origin string @@ -30,8 +28,6 @@ func Sink(ctx context.Context, c *cli.Config) *cobra.Command { Use: "signature", Short: "Creates signature for sink HTTP endpoint", RunE: func(cmd *cobra.Command, args []string) error { - c.InitServices(ctx, c) - method = strings.ToUpper(method) if expires != "" { diff --git a/system/commands/users.go b/system/commands/users.go index d58faf851..7eec1a934 100644 --- a/system/commands/users.go +++ b/system/commands/users.go @@ -1,7 +1,6 @@ package commands import ( - "context" "fmt" "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/pkg/rh" @@ -12,13 +11,14 @@ import ( "github.com/titpetric/factory" "golang.org/x/crypto/ssh/terminal" + "github.com/cortezaproject/corteza-server/pkg/auth" "github.com/cortezaproject/corteza-server/pkg/cli" "github.com/cortezaproject/corteza-server/system/repository" "github.com/cortezaproject/corteza-server/system/service" "github.com/cortezaproject/corteza-server/system/types" ) -func Users(ctx context.Context, c *cli.Config) *cobra.Command { +func Users() *cobra.Command { var ( flagNoPassword bool ) @@ -34,11 +34,9 @@ func Users(ctx context.Context, c *cli.Config) *cobra.Command { Use: "list", Short: "List users", Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( - db = factory.Database.MustGet("system") + ctx = auth.SetSuperUserContext(cli.Context()) + db = factory.Database.MustGet("system", "default") queryFlag = cmd.Flags().Lookup("query").Value.String() limitFlag = cmd.Flags().Lookup("limit").Value.String() @@ -95,11 +93,10 @@ func Users(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Add new user", Args: cobra.MinimumNArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( - db = factory.Database.MustGet("system") + ctx = auth.SetSuperUserContext(cli.Context()) + + db = factory.Database.MustGet("system", "default") userRepo = repository.User(ctx, db) authSvc = service.Auth(ctx) @@ -151,11 +148,9 @@ func Users(ctx context.Context, c *cli.Config) *cobra.Command { Short: "Change password for user", Args: cobra.MinimumNArgs(1), Run: func(cmd *cobra.Command, args []string) { - c.InitServices(ctx, c) - ctx = auth.SetSuperUserContext(ctx) - var ( - db = factory.Database.MustGet("system") + ctx = auth.SetSuperUserContext(cli.Context()) + db = factory.Database.MustGet("system", "default") userRepo = repository.User(ctx, db) authSvc = service.Auth(ctx) diff --git a/system/provision.go b/system/provision.go new file mode 100644 index 000000000..d4c3d41d8 --- /dev/null +++ b/system/provision.go @@ -0,0 +1,166 @@ +package system + +import ( + "context" + "io" + + "github.com/pkg/errors" + "github.com/titpetric/factory" + "go.uber.org/zap" + "gopkg.in/yaml.v2" + + "github.com/cortezaproject/corteza-server/pkg/auth" + impAux "github.com/cortezaproject/corteza-server/pkg/importer" + "github.com/cortezaproject/corteza-server/pkg/permissions" + "github.com/cortezaproject/corteza-server/pkg/settings" + provision "github.com/cortezaproject/corteza-server/provision/system" + "github.com/cortezaproject/corteza-server/system/importer" + "github.com/cortezaproject/corteza-server/system/repository" + "github.com/cortezaproject/corteza-server/system/service" + "github.com/cortezaproject/corteza-server/system/types" +) + +func provisionConfig(ctx context.Context, log *zap.Logger) (err error) { + log.Debug("running configuration provision") + + var provisioned bool + + // Make sure we have all full access for provisioning + ctx = auth.SetSuperUserContext(ctx) + + // if system is already provisioned, we do partial provisioning: + // missing settings only + if provisioned, err = isProvisioned(ctx); err != nil { + return err + } else if provisioned { + log.Debug("configuration already provisioned") + } + + readers, err := impAux.ReadStatic(provision.Asset) + if err != nil { + return err + } + + if provisioned { + return partialImportSettings(ctx, service.DefaultSettings, readers...) + } + + return errors.Wrap( + importer.Import(ctx, readers...), + "could not provision configuration for system service", + ) +} + +// Provision ONLY when there are no rules for role admins +func isProvisioned(ctx context.Context) (bool, error) { + return len(service.DefaultPermissions.FindRulesByRoleID(permissions.AdminsRoleID)) > 0, nil +} + +func makeDefaultApplications(ctx context.Context, log *zap.Logger) error { + db := factory.Database.MustGet("system", "default") + + repo := repository.Application(ctx, db) + + aa, _, err := repo.Find(types.ApplicationFilter{}) + if err != nil { + return err + } + + // List of apps to create. + // + // We use Unify.Url field for matching, + // so make sure it's always present! + defApps := types.ApplicationSet{ + &types.Application{ + Name: "CRM", + Enabled: true, + Unify: &types.ApplicationUnify{ + Name: "CRM", + Listed: true, + Icon: "/applications/crust_favicon.png", + Logo: "/applications/crust.jpg", + Url: "/compose/ns/crm/pages", + }, + }, + + &types.Application{ + Name: "Service Cloud", + Enabled: true, + Unify: &types.ApplicationUnify{ + Name: "Service Cloud", + Listed: true, + Icon: "/applications/crust_favicon.png", + Logo: "/applications/crust.jpg", + Url: "/compose/ns/service-cloud/pages", + }, + }, + } + + return defApps.Walk(func(defApp *types.Application) error { + for _, a := range aa { + if a.Unify != nil && a.Unify.Url == defApp.Unify.Url { + // App already added. + return nil + } + } + + defApp, err = repo.Create(defApp) + log.Info( + "creating default application", + zap.String("name", defApp.Name), + zap.Uint64("name", defApp.ID), + zap.Error(err), + ) + return err + + return nil + }) +} + +// Partial import of settings from provision files +func partialImportSettings(ctx context.Context, ss settings.Service, ff ...io.Reader) (err error) { + var ( + // decoded content from YAML files + aux interface{} + + si = settings.NewImporter() + + // importer w/o permissions & roles + // we need only settings + imp = importer.NewImporter(nil, si, nil) + + // current value + current settings.ValueSet + + // unexisting values + unex settings.ValueSet + ) + + for _, f := range ff { + if err = yaml.NewDecoder(f).Decode(&aux); err != nil { + return + } + + err = imp.Cast(aux) + if err != nil { + return + } + } + + // Get all "current" settings storage + current, err = ss.FindByPrefix(ctx) + if err != nil { + return + } + + // Compare current settings with imported, get all that do not exist yet + if unex = si.GetValues(); len(unex) > 0 { + // Store non existing + err = ss.BulkSet(ctx, current.New(unex)) + if err != nil { + return + } + } + + return nil +} diff --git a/system/provision_oidc.go b/system/provision_oidc.go new file mode 100644 index 000000000..8b37fba6f --- /dev/null +++ b/system/provision_oidc.go @@ -0,0 +1,124 @@ +package system + +import ( + "context" + "strings" + + "github.com/pkg/errors" + "go.uber.org/zap" + + "github.com/cortezaproject/corteza-server/pkg/app/options" + "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/system/auth/external" + "github.com/cortezaproject/corteza-server/system/types" +) + +// Provisions OIDC providers from PROVISION_OIDC_PROVIDER env variable +// +// Env variable should contains space delimited pairs of providers ( ....) +func oidcAutoDiscovery(ctx context.Context, log *zap.Logger) (err error) { + var provider = strings.TrimSpace(options.EnvString("", "PROVISION_OIDC_PROVIDER", "")) + + log.Debug("OIDC auto discovery provision", + zap.String("envkey", "PROVISION_OIDC_PROVIDER"), + zap.String("providers", provider), + ) + + if len(provider) == 0 { + return + } + + var ( + providers = strings.Split(provider, " ") + plen = len(providers) + name, purl string + eap *types.ExternalAuthProvider + ) + + if plen%2 == 1 { + return errors.New("expecting even number of providers") + } + + for p := 0; p < plen; p = p + 2 { + name, purl = providers[p], providers[p+1] + + // force: false + // because we do not want to override the provider every time the system restarts + // + // validate: false + // because at the initial (empty db) start, we can not actually validate (server is not yet up) + // + // enable: true + // we want provider & the entire external auth to be validated + eap, err = external.RegisterOidcProvider(ctx, name, purl, false, false, true) + + if err != nil { + log.Error( + "could not register OIDC provider", + zap.String("url", purl), + zap.String("name", name), + zap.Error(err)) + return + } else if eap == nil { + log.Info("provider already exists", + zap.String("name", name)) + } else { + log.Info("provider successfully registered", + zap.String("url", purl), + zap.String("key", eap.Key), + zap.String("name", name)) + } + } + + return +} + +func authAddExternals(ctx context.Context) (err error) { + var ( + kinds = []string{ + "github", + "facebook", + "google", + "linkedin", + "oidc", + } + + env, p, name string + + pp []string + + eap *types.ExternalAuthProvider + ) + + for _, kind := range kinds { + env = "PROVISION_SETTINGS_AUTH_EXTERNAL_" + strings.ToUpper(kind) + + p = strings.TrimSpace(options.EnvString("", env, "")) + if len(p) == 0 { + continue + } + + eap = &types.ExternalAuthProvider{Enabled: true} + + if kind == "oidc" { + pp = strings.SplitN(p, " ", 4) + + // Spread name, issuer-url, key and secret from provision string for OIDC provider + name, eap.IssuerUrl, eap.Key, eap.Secret = pp[0], pp[1], pp[2], pp[3] + + eap.Handle = external.OIDC_PROVIDER_PREFIX + name + } else { + pp = strings.SplitN(p, " ", 2) + + // Spread key and secret from provision string + eap.Key, eap.Secret = pp[0], pp[1] + eap.Handle = kind + } + + ctx = auth.SetSuperUserContext(ctx) + + _ = external.AddProvider(ctx, eap, false) + } + + return +} diff --git a/system/provision_settings.go b/system/provision_settings.go new file mode 100644 index 000000000..9633a9b96 --- /dev/null +++ b/system/provision_settings.go @@ -0,0 +1,368 @@ +package system + +import ( + "context" + "fmt" + "net/url" + "os" + "strings" + + "github.com/spf13/cast" + "go.uber.org/zap" + + "github.com/cortezaproject/corteza-server/pkg/logger" + "github.com/cortezaproject/corteza-server/pkg/rand" + settings "github.com/cortezaproject/corteza-server/pkg/settings" + "github.com/cortezaproject/corteza-server/system/service" +) + +var ( + IsMonolith bool +) + +// Discovers "auth.%" settings from the environment +// +// when other kinds of auto-discoverable settings come, lambdas inside will probably need a bit of refactoring +func authSettingsAutoDiscovery(ctx context.Context, log *zap.Logger, svc settings.Service) (err error) { + type ( + stringWrapper func() string + boolWrapper func() bool + ) + + var ( + current settings.ValueSet + ) + + if log == nil { + log = zap.NewNop() + } + + log = service.DefaultLogger.Named("auth-settings-discovery") + + current, err = svc.FindByPrefix(ctx, "auth.") + if err != nil { + return + } + + var ( + new = current + + // Setter + // + // Finds existing settings, tries with environmental "PROVISION_SETTINGS_AUTH_..." probing + // and falls back to default value + // + // We are extremely verbose here - we want to show all the info available and + // how settings were discovered and set + // + // @todo generalize and move under settings + set = func(name string, env string, def interface{}, maskSensitive bool) { + var ( + log = log.With( + zap.String("name", name), + ) + + v = current.First(name) + value interface{} + ) + + if v != nil { + // Nothing to discover, already set + log.Info("already set", logger.MaskIf("value", v, maskSensitive)) + return + } + + v = &settings.Value{Name: name} + + value, envExists := os.LookupEnv(env) + + switch dfn := def.(type) { + case stringWrapper: + log = log.With(zap.String("type", "string")) + // already a string, no need to do any magic + if envExists { + log = log.With(zap.String("env", env), logger.MaskIf("value", value, maskSensitive)) + } else { + value = dfn() + log = log.With(zap.Any("default", value)) + } + case boolWrapper: + log = log.With(zap.String("type", "bool")) + + if envExists { + value = cast.ToBool(value) + log = log.With(zap.String("env", env), zap.Any("value", value)) + } else { + value = dfn() + log = log.With(zap.Any("default", value)) + } + + default: + log.Error("unsupported type") + return + } + + if err := v.SetValue(value); err != nil { + log.Error("could not set value", zap.Error(err)) + return + } + + log.Info("value auto-discovered") + + new.Replace(v) + } + + // Default value functions + // + // all are wrapped (stringWrapper, boolWrapper) to delay execution + // of the function to the very last point + + frontendUrl = func(path string) stringWrapper { + const ( + feBase = "auth.frontend.url.base" + extRedir = "auth.external.redirect-url" + ) + + return func() (base string) { + base = new.First(feBase).String() + + if len(base) == 0 { + + // Not found, try to get it from the external redirect URL + redirURL := new.First(extRedir).String() + if len(redirURL) == 0 { + return + } + + log.Info( + "discovering frontend url from '"+extRedir+"'", + zap.String(extRedir, redirURL)) + + // Removing placeholder + redirURL = fmt.Sprintf(redirURL, "") + + p, err := url.Parse(redirURL) + if err != nil { + log.Error("could not parse '"+extRedir+"'", zap.Error(err)) + return + } + + h := p.Host + s := "api." + if i := strings.Index(h, s); i > 0 { + // If there is a "api." prefix in the hostname of the external redirect-uri value + // cut it off and use that as a frontend url base + h = h[i+len(s):] + } + + base = p.Scheme + "://" + h + } + + if len(base) > 0 { + return strings.TrimRight(base, "/") + path + } + + return "" + } + } + + // Assuming secure backend when redirect URL starts with https:// + isSecure = func() boolWrapper { + return func() bool { + return strings.Index(new.First("auth.external.redirect-url").String(), "https://") == 0 + } + } + + // Assume we have emailing capabilities if SMTP_HOST variable is set + emailCapabilities = func() boolWrapper { + return func() bool { + val, has := os.LookupEnv("SMTP_HOST") + return has && len(val) > 0 + } + } + + // Where should external authentication providers redirect to? + // we need to set full, absolute URL to the callback endpoint + externalAuthRedirectUrl = func() stringWrapper { + return func() string { + var ( + path = "/auth/external/%s/callback" + + // All env keys we'll check, first that has any value set, will be used as hostname + keysWithHostnames = []string{ + "DOMAIN", + "LETSENCRYPT_HOST", + "VIRTUAL_HOST", + "HOSTNAME", + "HOST", + } + ) + + // Prefix path if we're running wrapped as a monolith: + if IsMonolith { + path = "/system" + path + } + + log.Info("scanning env variables for hostname", zap.Strings("candidates", keysWithHostnames)) + + for _, key := range keysWithHostnames { + if host, has := os.LookupEnv(key); has { + log.Info("hostname env variable found", zap.String("env", key)) + // Make life easier for development in local environment, + // and set HTTP schema. Might cause problems if someone + // is using valid external hostname + if strings.Contains(host, "local.") { + return "http://" + host + path + } else { + return "https://" + host + path + } + } else { + } + } + + // Fallback is empty string + // this will cause error when doing OIDC auto-discovery (and we want that) + // @todo ^^ + return "" + } + } + + rand stringWrapper = func() string { + return string(rand.Bytes(64)) + } + + wrapBool = func(val bool) boolWrapper { + return func() bool { return val } + } + + wrapString = func(val string) stringWrapper { + return func() string { return val } + } + ) + + // List of name-value pairs we need to iterate and set + list := []struct { + // Setting name + nme string + + // provision environmental variable name + // we're using full variable name here so developers + // can find where things are comming from + env string + + // default value + // expects one of the *wrapper() functions + // this also determinate the value type of the setting and casting rules for the env value + def interface{} + + // mask value if sensitive + mask bool + }{ + // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // + // External auth + + // Enable external auth + { + "auth.external.enabled", + "PROVISION_SETTINGS_AUTH_EXTERNAL_ENABLED", + wrapBool(true), + false}, + + { + "auth.external.redirect-url", + "PROVISION_SETTINGS_AUTH_EXTERNAL_REDIRECT_URL", + externalAuthRedirectUrl(), + false}, + + { + "auth.external.session-store-secret", + "PROVISION_SETTINGS_AUTH_EXTERNAL_SESSION_STORE_SECRET", + rand, + true}, + + // Disable external auth + { + "auth.external.session-store-secure", + "PROVISION_SETTINGS_AUTH_EXTERNAL_SESSION_STORE_SECURE", + isSecure(), + false}, + + // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // + // Auth frontend + + { + "auth.frontend.url.base", + "PROVISION_SETTINGS_AUTH_FRONTEND_URL_BASE", + frontendUrl("/"), + false}, + + // @todo w/o token= + { + "auth.frontend.url.password-reset", + "PROVISION_SETTINGS_AUTH_FRONTEND_URL_PASSWORD_RESET", + frontendUrl("/auth/reset-password?token="), + false}, + + // @todo w/o token= + { + "auth.frontend.url.email-confirmation", + "PROVISION_SETTINGS_AUTH_FRONTEND_URL_EMAIL_CONFIRMATION", + frontendUrl("/auth/confirm-email?token="), + false}, + + // @todo check if this is correct?! + { + "auth.frontend.url.redirect", + "PROVISION_SETTINGS_AUTH_FRONTEND_URL_REDIRECT", + frontendUrl("/auth"), + false}, + + // Auth email + { + "auth.mail.from-address", + "PROVISION_SETTINGS_AUTH_EMAIL_FROM_ADDRESS", + wrapString("info@example.tld"), + false}, + + { + "auth.mail.from-name", + "PROVISION_SETTINGS_AUTH_EMAIL_FROM_NAME", + wrapString("Example Sender"), + false}, + + // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // // + // Enable internal login + { + "auth.internal.enabled", + "PROVISION_SETTINGS_AUTH_INTERNAL_ENABLED", + wrapBool(true), + false}, + + // Enable internal signup + { + "auth.internal.signup.enabled", + "PROVISION_SETTINGS_AUTH_INTERNAL_SIGNUP_ENABLED", + wrapBool(true), + false}, + + // Enable email confirmation if we have email capabilities + { + "auth.internal.signup.email-confirmation-required", + "PROVISION_SETTINGS_AUTH_INTERNAL_SIGNUP_EMAIL_CONFIRMATION_REQUIRED", + emailCapabilities(), + false}, + + // Enable password reset if we have email capabilities + { + "auth.internal.password-reset.enabled", + "PROVISION_SETTINGS_AUTH_INTERNAL_PASSWORD_RESET_ENABLED", + emailCapabilities(), + false}, + } + + for _, item := range list { + set(item.nme, item.env, item.def, item.mask) + } + + return svc.BulkSet(ctx, new) +} diff --git a/system/repository/repository.go b/system/repository/repository.go index 5383b6bec..1b7eb2dca 100644 --- a/system/repository/repository.go +++ b/system/repository/repository.go @@ -15,7 +15,7 @@ type ( // DB produces a contextual DB handle func DB(ctx context.Context) *factory.DB { - return factory.Database.MustGet("system").With(ctx) + return factory.Database.MustGet("system", "default").With(ctx) } // With updates repository and database contexts diff --git a/system/rest/stats.go b/system/rest/stats.go index 04bcdd913..906e78419 100644 --- a/system/rest/stats.go +++ b/system/rest/stats.go @@ -23,7 +23,7 @@ type ( func (Stats) New() *Stats { return &Stats{ - svc: service.Statistics(context.Background()), + svc: service.DefaultStatistics, } } diff --git a/system/service/service.go b/system/service/service.go index 2dad8b0ba..ae0dcb4d9 100644 --- a/system/service/service.go +++ b/system/service/service.go @@ -92,9 +92,11 @@ var ( DefaultOrganisation OrganisationService DefaultApplication ApplicationService DefaultReminder ReminderService + + DefaultStatistics *statistics ) -func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { +func Initialize(ctx context.Context, log *zap.Logger, c Config) (err error) { DefaultLogger = log.Named("service") if DefaultPermissions == nil { @@ -112,12 +114,6 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { CurrentSettings, ) - // Run initial update of current settings with super-user credentials - err = DefaultSettings.UpdateCurrent(intAuth.SetSuperUserContext(ctx)) - if err != nil { - return - } - DefaultUser = User(ctx) DefaultRole = Role(ctx) DefaultOrganisation = Organisation(ctx) @@ -178,6 +174,17 @@ func Init(ctx context.Context, log *zap.Logger, c Config) (err error) { } DefaultSink = Sink() + DefaultStatistics = Statistics(ctx) + + return +} + +func Activate(ctx context.Context) (err error) { + // Run initial update of current settings with super-user credentials + err = DefaultSettings.UpdateCurrent(intAuth.SetSuperUserContext(ctx)) + if err != nil { + return + } return } diff --git a/system/system.go b/system/system.go deleted file mode 100644 index 67711667b..000000000 --- a/system/system.go +++ /dev/null @@ -1,170 +0,0 @@ -package system - -import ( - "context" - "net" - - _ "github.com/joho/godotenv/autoload" - "github.com/spf13/cobra" - "github.com/titpetric/factory" - "go.uber.org/zap" - - "github.com/cortezaproject/corteza-server/pkg/auth" - "github.com/cortezaproject/corteza-server/pkg/cli" - "github.com/cortezaproject/corteza-server/system/auth/external" - "github.com/cortezaproject/corteza-server/system/commands" - migrate "github.com/cortezaproject/corteza-server/system/db" - "github.com/cortezaproject/corteza-server/system/grpc" - "github.com/cortezaproject/corteza-server/system/rest" - "github.com/cortezaproject/corteza-server/system/service" -) - -const ( - system = "system" -) - -func Configure() *cli.Config { - var ( - servicesInitialized bool - ) - - return &cli.Config{ - ServiceName: system, - - RootCommandPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) (err error) { - return - }, - }, - - InitServices: func(ctx context.Context, c *cli.Config) { - if servicesInitialized { - return - } - servicesInitialized = true - - // storagePath := options.EnvString("", "SYSTEM_STORAGE_PATH", "var/store") - cli.HandleError(service.Init(ctx, c.Log, service.Config{ - Corredor: *c.ScriptRunner, - })) - - }, - - ApiServerPreRun: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - if c.ProvisionOpt.MigrateDatabase { - cli.HandleError(c.ProvisionMigrateDatabase.Run(ctx, cmd, c)) - } - - c.InitServices(ctx, c) - - if c.ProvisionOpt.Configuration { - ctx = auth.SetSuperUserContext(ctx) - - // read system's config files (YAML) - cli.HandleError(provisionConfig(ctx, cmd, c)) - - // creates default applications that will appear in Unify/One - // @todo migrate this to provisioning/YAML - cli.HandleError(makeDefaultApplications(ctx, cmd, c)) - - // auto-discovery auth.* settings - cli.HandleError(authSettingsAutoDiscovery(ctx, c.Log, service.DefaultSettings)) - - // external provider auto configuration - // creates: auth.external.providers.(google|linkedin|github|facebook).* - cli.HandleError(authAddExternals(ctx, cmd, c)) - - // OIDC provider auto configuration - // creates: auth.external.providers.openid-connect.* - cli.HandleError(oidcAutoDiscovery(ctx, cmd, c)) - } - - { - var ( - grpcLog = c.Log.Named("grpc-server") - grpcLogConn = grpcLog.With(zap.String("addr", c.GRPCServerSystem.Addr)) - ) - - // Temporary gRPC server initialization location - // @todo move out of system Configure - grpcServer := grpc.NewServer() - - ln, err := net.Listen(c.GRPCServerSystem.Network, c.GRPCServerSystem.Addr) - if err != nil { - grpcLogConn.Error("could not start gRPC server", zap.Error(err)) - } - - go func() { - select { - case <-ctx.Done(): - grpcLogConn.Debug("shutting down") - grpcServer.GracefulStop() - _ = ln.Close() - } - }() - - go func() { - grpcLogConn.Info("Starting gRPC server") - err := grpcServer.Serve(ln) - grpcLogConn.Info("stopped", zap.Error(err)) - }() - } - - // Initialize external authentication (from default settings) - external.Init() - go service.Watchers(ctx) - return nil - }, - }, - - ApiServerRoutes: cli.Mounters{ - rest.MountRoutes, - }, - - AdtSubCommands: cli.CommandMakers{ - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Settings(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Auth(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Importer(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Exporter(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Users(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Roles(ctx, c) - }, - func(ctx context.Context, c *cli.Config) *cobra.Command { - return commands.Sink(ctx, c) - }, - }, - - ProvisionMigrateDatabase: cli.Runners{ - func(ctx context.Context, cmd *cobra.Command, c *cli.Config) error { - if !c.ProvisionOpt.MigrateDatabase { - return nil - } - - var db, err = factory.Database.Get(system) - if err != nil { - return err - } - - db = db.With(ctx).Quiet() - - return migrate.Migrate(db, c.Log) - }, - }, - - ProvisionConfig: cli.Runners{ - provisionConfig, - }, - } -} diff --git a/system/system_test.go b/system/system_test.go deleted file mode 100644 index 891f02b18..000000000 --- a/system/system_test.go +++ /dev/null @@ -1,15 +0,0 @@ -package system - -import ( - "context" - "testing" - - "github.com/stretchr/testify/require" -) - -func TestConfigure(t *testing.T) { - var config = Configure() - require.True(t, config != nil, "Configure valid") - require.True(t, func() bool { config.Init(); return true }(), "Initialization ok") - require.True(t, config.MakeCLI(context.Background()) != nil, "CLI created") -} diff --git a/tests/compose/main_test.go b/tests/compose/main_test.go index 426bba95a..77790f183 100644 --- a/tests/compose/main_test.go +++ b/tests/compose/main_test.go @@ -15,7 +15,6 @@ import ( "go.uber.org/zap" "github.com/cortezaproject/corteza-server/compose" - migrate "github.com/cortezaproject/corteza-server/compose/db" "github.com/cortezaproject/corteza-server/compose/rest" "github.com/cortezaproject/corteza-server/compose/service" "github.com/cortezaproject/corteza-server/compose/types" @@ -41,13 +40,14 @@ type ( ) var ( - cfg *cli.Config - r chi.Router - p = &permissions.TestService{} + inited bool + app = &compose.App{} + r chi.Router + p = &permissions.TestService{} ) func db() *factory.DB { - return factory.Database.MustGet("compose").With(context.Background()) + return factory.Database.MustGet("compose", "default").With(context.Background()) } // random string, 10 chars long by default @@ -63,7 +63,7 @@ func rs(a ...int) string { func InitConfig() { var err error - if cfg != nil { + if inited { return } @@ -71,21 +71,13 @@ func InitConfig() { ctx := context.Background() log, _ := zap.NewDevelopment() + logger.SetDefault(log) - cfg = compose.Configure() - cfg.Log = log + cli.HandleError(app.Connect(ctx)) + auth.SetupDefault(rs(32), 10) + cli.HandleError(app.ProvisionMigrateDatabase(ctx)) - cfg.Init() - - auth.SetupDefault(string(rand.Bytes(32)), 10) - - if err = cfg.RootCommandDBSetup.Run(ctx, nil, cfg); err != nil { - panic(err) - } else if err := migrate.Migrate(factory.Database.MustGet("compose"), log); err != nil { - panic(err) - } - - p = permissions.NewTestService(ctx, cfg.Log, db(), "compose_permission_rules") + p = permissions.NewTestService(ctx, log, db(), "compose_permission_rules") logger.SetDefault(log) service.DefaultPermissions = p @@ -93,7 +85,8 @@ func InitConfig() { panic(err) } - cfg.InitServices(ctx, cfg) + cli.HandleError(app.Initialize(ctx)) + inited = true } func InitApp() { @@ -105,7 +98,7 @@ func InitApp() { } r = chi.NewRouter() - r.Use(api.Base(logger.Default())...) + r.Use(api.BaseMiddleware(logger.Default())...) helpers.BindAuthMiddleware(r) rest.MountRoutes(r) } @@ -146,6 +139,7 @@ func (h helper) apiInit() *apitest.APITest { New(). Handler(r). Intercept(helpers.ReqHeaderAuthBearer(h.cUser)) + } func (h helper) mockPermissions(rules ...*permissions.Rule) { diff --git a/tests/messaging/main_test.go b/tests/messaging/main_test.go index d692bed77..1e32ab4bd 100644 --- a/tests/messaging/main_test.go +++ b/tests/messaging/main_test.go @@ -15,7 +15,6 @@ import ( "go.uber.org/zap" "github.com/cortezaproject/corteza-server/messaging" - migrate "github.com/cortezaproject/corteza-server/messaging/db" "github.com/cortezaproject/corteza-server/messaging/rest" "github.com/cortezaproject/corteza-server/messaging/service" "github.com/cortezaproject/corteza-server/messaging/types" @@ -41,19 +40,30 @@ type ( ) var ( - cfg *cli.Config - r chi.Router - p = &permissions.TestService{} + inited bool + app = &messaging.App{} + r chi.Router + p = &permissions.TestService{} ) func db() *factory.DB { - return factory.Database.MustGet("messaging").With(context.Background()) + return factory.Database.MustGet("messaging", "default").With(context.Background()) +} + +// random string, 10 chars long by default +func rs(a ...int) string { + var l = 10 + if len(a) > 0 { + l = a[0] + } + + return string(rand.Bytes(l)) } func InitConfig() { var err error - if cfg != nil { + if inited { return } @@ -61,21 +71,13 @@ func InitConfig() { ctx := context.Background() log, _ := zap.NewDevelopment() + logger.SetDefault(log) - cfg = messaging.Configure() - cfg.Log = log + cli.HandleError(app.Connect(ctx)) + auth.SetupDefault(rs(32), 10) + cli.HandleError(app.ProvisionMigrateDatabase(ctx)) - cfg.Init() - - auth.SetupDefault(string(rand.Bytes(32)), 10) - - if err = cfg.RootCommandDBSetup.Run(ctx, nil, cfg); err != nil { - panic(err) - } else if err := migrate.Migrate(factory.Database.MustGet("messaging"), log); err != nil { - panic(err) - } - - p = permissions.NewTestService(ctx, cfg.Log, db(), "messaging_permission_rules") + p = permissions.NewTestService(ctx, log, db(), "messaging_permission_rules") logger.SetDefault(log) service.DefaultPermissions = p @@ -83,7 +85,8 @@ func InitConfig() { panic(err) } - cfg.InitServices(ctx, cfg) + cli.HandleError(app.Initialize(ctx)) + inited = true } func InitApp() { @@ -95,7 +98,7 @@ func InitApp() { } r = chi.NewRouter() - r.Use(api.Base(logger.Default())...) + r.Use(api.BaseMiddleware(logger.Default())...) helpers.BindAuthMiddleware(r) rest.MountRoutes(r) } diff --git a/tests/system/main_test.go b/tests/system/main_test.go index c8fd649c6..8fa2005b9 100644 --- a/tests/system/main_test.go +++ b/tests/system/main_test.go @@ -20,7 +20,6 @@ import ( "github.com/cortezaproject/corteza-server/pkg/permissions" "github.com/cortezaproject/corteza-server/pkg/rand" "github.com/cortezaproject/corteza-server/system" - migrate "github.com/cortezaproject/corteza-server/system/db" "github.com/cortezaproject/corteza-server/system/rest" "github.com/cortezaproject/corteza-server/system/service" "github.com/cortezaproject/corteza-server/system/types" @@ -38,9 +37,10 @@ type ( ) var ( - cfg *cli.Config - r chi.Router - p = &permissions.TestService{} + inited bool + app = &system.App{} + r chi.Router + p = &permissions.TestService{} ) // random string, 10 chars long by default @@ -54,13 +54,11 @@ func rs(a ...int) string { } func db() *factory.DB { - return factory.Database.MustGet("system").With(context.Background()) + return factory.Database.MustGet("system", "default").With(context.Background()) } func InitConfig() { - var err error - - if cfg != nil { + if inited { return } @@ -68,29 +66,19 @@ func InitConfig() { ctx := context.Background() log, _ := zap.NewDevelopment() + logger.SetDefault(log) - cfg = system.Configure() - cfg.Log = log + cli.HandleError(app.Connect(ctx)) + auth.SetupDefault(rs(32), 10) + cli.HandleError(app.ProvisionMigrateDatabase(ctx)) - cfg.Init() - - auth.SetupDefault(string(rand.Bytes(32)), 10) - - if err = cfg.RootCommandDBSetup.Run(ctx, nil, cfg); err != nil { - panic(err) - } else if err := migrate.Migrate(factory.Database.MustGet("system"), log); err != nil { - panic(err) - } - - p = permissions.NewTestService(ctx, cfg.Log, db(), "sys_permission_rules") + p = permissions.NewTestService(ctx, log, db(), "sys_permission_rules") logger.SetDefault(log) service.DefaultPermissions = p - // if service.DefaultStore, err = store.NewWithAfero(afero.NewMemMapFs(), "test"); err != nil { - // panic(err) - // } - cfg.InitServices(ctx, cfg) + cli.HandleError(app.Initialize(ctx)) + inited = true } func InitApp() { @@ -102,7 +90,7 @@ func InitApp() { } r = chi.NewRouter() - r.Use(api.Base(logger.Default())...) + r.Use(api.BaseMiddleware(logger.Default())...) helpers.BindAuthMiddleware(r) rest.MountRoutes(r) }