3
0

Refactor bootstraping procedure

This commit is contained in:
Denis Arh
2020-08-15 18:28:58 +02:00
parent 1147830663
commit f4e6d2ae5a
53 changed files with 718 additions and 1462 deletions
+54
View File
@@ -0,0 +1,54 @@
package app
import (
"context"
"github.com/go-chi/chi"
"github.com/spf13/cobra"
"go.uber.org/zap"
"google.golang.org/grpc"
)
type (
httpApiServer interface {
MountRoutes(mm ...func(chi.Router))
Serve(ctx context.Context)
}
grpcServer interface {
RegisterServices(func(server *grpc.Server))
Serve(ctx context.Context)
}
CortezaApp struct {
Opt *Options
lvl int
Log *zap.Logger
// Store interface
//
// Just a blank interface{} because we want to avoid generating
// whole store interface (as we do for other packages).
//
// Value will be type-casted when assigned to sys/msg/cmp services
// with warnings when incompatible
Store interface{}
// CLI Commands
Command *cobra.Command
// Servers
HttpServer httpApiServer
GrpcServer grpcServer
}
)
func New() *CortezaApp {
app := &CortezaApp{
Opt: NewOptions(),
lvl: bootLevelWaiting,
}
app.InitCLI()
return app
}
+297
View File
@@ -0,0 +1,297 @@
package app
import (
"context"
"errors"
"fmt"
cmpService "github.com/cortezaproject/corteza-server/compose/service"
cmpEvent "github.com/cortezaproject/corteza-server/compose/service/event"
msgService "github.com/cortezaproject/corteza-server/messaging/service"
msgEvent "github.com/cortezaproject/corteza-server/messaging/service/event"
"github.com/cortezaproject/corteza-server/pkg/actionlog"
"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/eventbus"
"github.com/cortezaproject/corteza-server/pkg/healthcheck"
"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/scheduler"
"github.com/cortezaproject/corteza-server/pkg/sentry"
"github.com/cortezaproject/corteza-server/provision/compose"
"github.com/cortezaproject/corteza-server/provision/messaging"
"github.com/cortezaproject/corteza-server/provision/system"
"github.com/cortezaproject/corteza-server/store/mysql"
"github.com/cortezaproject/corteza-server/system/auth/external"
sysService "github.com/cortezaproject/corteza-server/system/service"
sysEvent "github.com/cortezaproject/corteza-server/system/service/event"
"go.uber.org/zap"
"time"
)
type (
storeUpgrader interface {
Upgrade(context.Context, *zap.Logger) error
}
)
const (
bootLevelWaiting = iota
bootLevelSetup
bootLevelStoreInitialized
bootLevelServicesInitialized
bootLevelUpgraded
bootLevelProvisioned
bootLevelActivated
)
// Setup configures all required services
func (app *CortezaApp) Setup() (err error) {
app.Log = logger.Default()
if app.lvl >= bootLevelSetup {
// Are basics already set-up?
return nil
}
hcd := healthcheck.Defaults()
hcd.Add(scheduler.Healthcheck, "Scheduler")
hcd.Add(mail.Healthcheck, "Mail")
hcd.Add(corredor.Healthcheck, "Corredor")
hcd.Add(db.Healthcheck(), "Database")
if err = sentry.Init(app.Opt.Sentry); err != nil {
return fmt.Errorf("could not initialize Sentry: %w", err)
}
// Use Sentry right away to handle any panics
// that might occur inside auth, mail setup...
defer sentry.Recover()
auth.SetupDefault(app.Opt.Auth.Secret, int(app.Opt.Auth.Expiry/time.Minute))
mail.SetupDialer(app.Opt.SMTP.Host, app.Opt.SMTP.Port, app.Opt.SMTP.User, app.Opt.SMTP.Pass, app.Opt.SMTP.From)
http.SetupDefaults(
app.Opt.HTTPClient.HttpClientTimeout,
app.Opt.HTTPClient.ClientTSLInsecure,
)
monitor.Setup(app.Log, app.Opt.Monitor)
scheduler.Setup(app.Log, eventbus.Service(), 0)
scheduler.Service().OnTick(
sysEvent.SystemOnInterval(),
sysEvent.SystemOnTimestamp(),
cmpEvent.ComposeOnInterval(),
cmpEvent.ComposeOnTimestamp(),
msgEvent.MessagingOnInterval(),
msgEvent.MessagingOnTimestamp(),
)
if err = corredor.Setup(app.Log, app.Opt.Corredor); err != nil {
return err
}
app.lvl = bootLevelSetup
return
}
// InitStore initializes store backend(s) and runs upgrade procedures
func (app *CortezaApp) InitStore(ctx context.Context) (err error) {
if app.lvl >= bootLevelStoreInitialized {
// Is store already initialised?
return nil
} else if err = app.Setup(); err != nil {
// Initialize previous level
return err
}
defer sentry.Recover()
// @todo this should be configurable
app.Store, err = mysql.New(ctx, app.Opt.DB.DSN)
if err != nil {
return err
}
if upgradableStore, ok := app.Store.(storeUpgrader); !ok {
app.Log.Debug("store does not support upgrades")
} else if !app.Opt.Upgrade.Always {
app.Log.Debug("store upgrade skipped (UPGRADE_ALWAYS=false)")
} else {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Upgrade)
// If not explicitly set (UPGRADE_DEBUG=true) suppress logging in upgrader
log := zap.NewNop()
if app.Opt.Upgrade.Debug {
log = app.Log.Named("store.upgrade")
log.Debug("store upgrade running in debug mode (UPGRADE_DEBUG=true)")
} else {
app.Log.Debug("store upgrade running (to enable upgrade debug logging set UPGRADE_DEBUG=true)")
}
if err = upgradableStore.Upgrade(ctx, log); err != nil {
return err
}
}
// deprecated connector
// current state of Corteza (repos...) still requires it
_, err = db.TryToConnect(ctx, app.Log, app.Opt.DB)
if err != nil {
return fmt.Errorf("could not connect to database: %w", err)
}
app.lvl = bootLevelStoreInitialized
return nil
}
// InitServices initializes all services used
func (app *CortezaApp) InitServices(ctx context.Context) (err error) {
if app.lvl >= bootLevelServicesInitialized {
return nil
} else if err := app.InitStore(ctx); err != nil {
return err
}
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Init)
defer sentry.Recover()
if err = corredor.Service().Connect(ctx); err != nil {
return
}
corredor.Service().SetUserFinder(sysService.DefaultUser)
corredor.Service().SetRoleFinder(sysService.DefaultRole)
// Initializes system services
//
// Note: this is a legacy approach, all services from all 3 apps
// will most likely be merged in the future
err = sysService.Initialize(ctx, app.Log, app.Store, sysService.Config{
ActionLog: app.Opt.ActionLog,
Storage: app.Opt.Storage,
})
if err != nil {
return
}
// Initializes compose services
//
// Note: this is a legacy approach, all services from all 3 apps
// will most likely be merged in the future
err = cmpService.Initialize(ctx, app.Log, app.Store, cmpService.Config{
ActionLog: app.Opt.ActionLog,
Storage: app.Opt.Storage,
})
if err != nil {
return
}
// Initializes messaging services
//
// Note: this is a legacy approach, all services from all 3 apps
// will most likely be merged in the future
err = msgService.Initialize(ctx, app.Log, app.Store, msgService.Config{
ActionLog: app.Opt.ActionLog,
Storage: app.Opt.Storage,
})
if err != nil {
return
}
// Initialize external authentication (from default settings)
external.Init()
app.lvl = bootLevelServicesInitialized
return
}
// Provision instance with configuration and settings
// by importing preset configurations and running autodiscovery procedures
func (app *CortezaApp) Provision(ctx context.Context) (err error) {
if app.lvl >= bootLevelProvisioned {
return
}
if err = app.InitServices(ctx); err != nil {
return err
}
if !app.Opt.Provision.Always {
app.Log.Debug("provisioning skipped (PROVISION_ALWAYS=false)")
} else {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Provision)
defer sentry.Recover()
ctx = auth.SetSuperUserContext(ctx)
if err = system.Provision(ctx, app.Log); err != nil {
return fmt.Errorf("could not provision messaging: %w", err)
}
if err = compose.Provision(ctx, app.Log); err != nil {
return fmt.Errorf("could not provision messaging: %w", err)
}
if err = messaging.Provision(ctx, app.Log); err != nil {
return fmt.Errorf("could not provision messaging: %w", err)
}
for errors.Unwrap(err) != nil {
err = errors.Unwrap(err)
}
if err != nil {
return err
}
}
app.lvl = bootLevelProvisioned
return
}
// Activate start all internal services and watchers
func (app *CortezaApp) Activate(ctx context.Context) (err error) {
if app.lvl >= bootLevelActivated {
return
} else if err := app.Provision(ctx); err != nil {
return err
}
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Activate)
defer sentry.Recover()
// Start scheduler
scheduler.Service().Start(ctx)
// Load corredor scripts & init watcher (script reloader)
corredor.Service().Load(ctx)
corredor.Service().Watch(ctx)
sysService.Watchers(ctx)
cmpService.Watchers(ctx)
msgService.Watchers(ctx)
if err = sysService.Activate(ctx); err != nil {
return err
}
if err = cmpService.Activate(ctx); err != nil {
return err
}
if err = msgService.Activate(ctx); err != nil {
return err
}
app.lvl = bootLevelActivated
return nil
}
+55
View File
@@ -0,0 +1,55 @@
package app
import (
"github.com/cortezaproject/corteza-server/pkg/cli"
systemCommands "github.com/cortezaproject/corteza-server/system/commands"
)
// CLI function initializes basic Corteza subsystems
// and sets-up the command line interface
func (app *CortezaApp) InitCLI() {
ctx := cli.Context()
app.Command = cli.RootCommand(nil)
serveCmd := cli.ServeCommand(func() (err error) {
if err = app.Activate(ctx); err != nil {
return
}
return app.Serve(ctx)
})
upgradeCmd := cli.UpgradeCommand(func() (err error) {
if err = app.InitStore(ctx); err != nil {
return
}
return
})
provisionCmd := cli.ProvisionCommand(func() (err error) {
if err = app.Provision(ctx); err != nil {
return
}
return
})
app.Command.AddCommand(
systemCommands.Users(app),
systemCommands.Roles(app),
systemCommands.Auth(app),
systemCommands.RBAC(app),
systemCommands.Sink(app),
serveCmd,
upgradeCmd,
provisionCmd,
cli.VersionCommand(),
)
}
func (app *CortezaApp) Execute() error {
return app.Command.Execute()
}
+1 -3
View File
@@ -1,7 +1,7 @@
package app
import (
"github.com/cortezaproject/corteza-server/pkg/app/options"
"github.com/cortezaproject/corteza-server/pkg/options"
)
type (
@@ -19,7 +19,6 @@ type (
Monitor options.MonitorOpt
WaitFor options.WaitForOpt
HTTPServer options.HTTPServerOpt
GRPCServer options.GRPCServerOpt
Websocket options.WebsocketOpt
}
)
@@ -44,7 +43,6 @@ func NewOptions(prefix ...string) *Options {
Monitor: *options.Monitor(p),
WaitFor: *options.WaitFor(p),
HTTPServer: *options.HTTP(p),
GRPCServer: *options.GRPCServer(p),
Websocket: *options.Websocket(p),
}
}
+62
View File
@@ -0,0 +1,62 @@
package app
import (
"context"
composeRest "github.com/cortezaproject/corteza-server/compose/rest"
messagingRest "github.com/cortezaproject/corteza-server/messaging/rest"
"github.com/cortezaproject/corteza-server/pkg/actionlog"
"github.com/cortezaproject/corteza-server/pkg/api"
"github.com/cortezaproject/corteza-server/pkg/webapp"
systemRest "github.com/cortezaproject/corteza-server/system/rest"
"github.com/go-chi/chi"
"strings"
"sync"
)
func (app *CortezaApp) Serve(ctx context.Context) (err error) {
wg := &sync.WaitGroup{}
{
// @todo refactor wait-for out of HTTP API server.
app.HttpServer = api.New(app.Log, app.Opt.HTTPServer, app.Opt.WaitFor)
app.HttpServer.MountRoutes(app.mountHttpRoutes)
wg.Add(1)
go func() {
app.HttpServer.Serve(actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_API_REST))
wg.Done()
}()
}
{
//wg.Add(1)
//go func(ctx context.Context) {
// grpcApi.Serve(actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_API_GRPC))
// wg.Done()
//}(ctx)
}
// Wait for all servers to be done
wg.Wait()
return nil
}
func (app *CortezaApp) mountHttpRoutes(r chi.Router) {
var (
apiBaseUrl = strings.Trim(app.Opt.HTTPServer.ApiBaseUrl, "/")
webappBaseUrl = strings.Trim(app.Opt.HTTPServer.WebappBaseUrl, "/")
)
if app.Opt.HTTPServer.ApiEnabled {
r.Route("/"+apiBaseUrl, func(r chi.Router) {
r.Route("/system", systemRest.MountRoutes)
r.Route("/compose", composeRest.MountRoutes)
r.Route("/messaging", messagingRest.MountRoutes)
})
}
if app.Opt.HTTPServer.WebappEnabled {
r.Route("/"+webappBaseUrl, webapp.MakeWebappServer(app.Opt.HTTPServer))
}
}
-19
View File
@@ -1,19 +0,0 @@
package main
import (
"github.com/cortezaproject/corteza-server/compose"
"github.com/cortezaproject/corteza-server/corteza"
"github.com/cortezaproject/corteza-server/pkg/app"
"github.com/cortezaproject/corteza-server/pkg/logger"
)
func main() {
logger.Init()
app.Run(
logger.Default(),
app.NewOptions(compose.SERVICE),
&corteza.App{},
&compose.App{},
)
}
+14
View File
@@ -0,0 +1,14 @@
package main
import (
"github.com/cortezaproject/corteza-server/app"
"github.com/cortezaproject/corteza-server/pkg/cli"
"github.com/cortezaproject/corteza-server/pkg/logger"
)
func main() {
// Initialize logger before any other action
logger.Init()
cli.HandleError(app.New().Execute())
}
-19
View File
@@ -1,19 +0,0 @@
package main
import (
"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/logger"
)
func main() {
logger.Init()
app.Run(
logger.Default(),
app.NewOptions(messaging.SERVICE),
&corteza.App{},
&messaging.App{},
)
}
-26
View File
@@ -1,26 +0,0 @@
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/app"
"github.com/cortezaproject/corteza-server/pkg/logger"
"github.com/cortezaproject/corteza-server/system"
)
func main() {
logger.Init()
app.Run(
logger.Default(),
app.NewOptions(),
&corteza.App{},
&monolith.App{
System: &system.App{},
Compose: &compose.App{},
Messaging: &messaging.App{},
},
)
}
-19
View File
@@ -1,19 +0,0 @@
package main
import (
"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() {
logger.Init()
app.Run(
logger.Default(),
app.NewOptions(system.SERVICE),
&corteza.App{},
&system.App{},
)
}
-110
View File
@@ -1,110 +0,0 @@
package compose
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/automation"
"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/compose/service/event"
"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/scheduler"
"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
scheduler.Service().OnTick(
event.ComposeOnInterval(),
event.ComposeOnTimestamp(),
)
// @todo Wire in cross-service JWT maker for Corredor
corredor.Service().SetUserFinder(nil)
corredor.Service().SetRoleFinder(nil)
return
}
func (app *App) Upgrade(ctx context.Context) (err error) {
db := factory.Database.MustGet().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{
ActionLog: app.Opts.ActionLog,
Storage: app.Opts.Storage,
})
if err != nil {
return
}
// Initialize external authenpkg/corredor/conn_test.go:62:12:tication (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)
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(),
commands.NGImporter(),
// temp command, will be removed in 2020.6
automation.ScriptExporter(SERVICE),
)
}
-11
View File
@@ -1,11 +0,0 @@
package compose
import (
"testing"
"github.com/cortezaproject/corteza-server/pkg/app"
)
func TestConfigure(t *testing.T) {
var _ app.Runnable = &App{}
}
-105
View File
@@ -1,105 +0,0 @@
package corteza
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/healthcheck"
"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/eventbus"
"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/scheduler"
"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) {
hcd := healthcheck.Defaults()
hcd.Add(scheduler.Healthcheck, "Scheduler")
hcd.Add(mail.Healthcheck, "Mail")
hcd.Add(corredor.Healthcheck, "Corredor")
hcd.Add(db.Healthcheck(), "Database")
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.Auth.Secret, int(opts.Auth.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)
scheduler.Setup(log, eventbus.Service(), 0)
if err = corredor.Setup(log, opts.Corredor); err != nil {
return err
}
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")
}
if err = corredor.Service().Connect(ctx); err != nil {
return
}
return
}
func (app *App) Upgrade(ctx context.Context) error {
return nil
}
func (app *App) Activate(ctx context.Context) (err error) {
// Start scheduler
scheduler.Service().Start(ctx)
// Load corredor scripts & init watcher (script reloader)
corredor.Service().Load(ctx)
corredor.Service().Watch(ctx)
return nil
}
func (app *App) Provision(ctx context.Context) error {
return nil
}
-112
View File
@@ -1,112 +0,0 @@
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/service/event"
"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"
"github.com/cortezaproject/corteza-server/pkg/scheduler"
)
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
scheduler.Service().OnTick(
event.MessagingOnInterval(),
event.MessagingOnTimestamp(),
)
app.ws = websocket.New(&websocket.Config{
Timeout: opts.Websocket.Timeout,
PingTimeout: opts.Websocket.PingTimeout,
PingPeriod: opts.Websocket.PingPeriod,
})
// @todo Wire in cross-service JWT maker for Corredor
corredor.Service().SetUserFinder(nil)
corredor.Service().SetRoleFinder(nil)
return
}
func (app *App) Upgrade(ctx context.Context) (err error) {
db := factory.Database.MustGet().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{
ActionLog: app.Opts.ActionLog,
Storage: app.Opts.Storage,
})
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)
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(),
)
}
-11
View File
@@ -1,11 +0,0 @@
package messaging
import (
"testing"
"github.com/cortezaproject/corteza-server/pkg/app"
)
func TestConfigure(t *testing.T) {
var _ app.Runnable = &App{}
}
-182
View File
@@ -1,182 +0,0 @@
package monolith
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/webapp"
"github.com/go-chi/chi"
_ "github.com/joho/godotenv/autoload"
"github.com/pkg/errors"
"github.com/spf13/cobra"
"go.uber.org/zap"
"google.golang.org/grpc"
"strings"
"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"
systemService "github.com/cortezaproject/corteza-server/system/service"
)
type (
App struct {
Opts *app.Options
Core *corteza.App
System *system.App
Compose *compose.App
Messaging *messaging.App
}
)
func (monolith *App) Setup(log *zap.Logger, opts *app.Options) (err error) {
monolith.Opts = opts
// Make sure system behaves properly
//
// This will alter the auth settings provision procedure
system.IsMonolith = true
if err = monolith.CheckOptions(); err != nil {
return
}
err = app.RunSetup(
log,
opts,
monolith.System,
monolith.Compose,
monolith.Messaging,
)
if err != nil {
return
}
return
}
func (monolith *App) CheckOptions() error {
o := monolith.Opts.HTTPServer
if o.ApiEnabled && o.WebappEnabled && o.ApiBaseUrl == o.WebappBaseUrl {
return errors.Errorf("cannot serve api and web apps form the same base url (%v)", o.ApiBaseUrl)
}
return nil
}
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) {
err = app.RunInitialize(
ctx,
monolith.System,
monolith.Compose,
monolith.Messaging,
)
if err != nil {
return
}
corredor.Service().SetUserFinder(systemService.DefaultUser)
corredor.Service().SetRoleFinder(systemService.DefaultRole)
return
}
func (monolith *App) Activate(ctx context.Context) (err error) {
err = app.RunActivate(
ctx,
monolith.System,
monolith.Compose,
monolith.Messaging,
)
if err != nil {
return
}
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) {
var (
apiEnabled = monolith.Opts.HTTPServer.ApiEnabled
apiBaseUrl = strings.Trim(monolith.Opts.HTTPServer.ApiBaseUrl, "/")
webappEnabled = monolith.Opts.HTTPServer.WebappEnabled
webappBaseUrl = strings.Trim(monolith.Opts.HTTPServer.WebappBaseUrl, "/")
)
if apiEnabled {
r.Route("/"+apiBaseUrl, func(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)
})
})
}
if webappEnabled {
r.Route("/"+webappBaseUrl, webapp.MakeWebappServer(monolith.Opts.HTTPServer))
}
}
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,
)
}
-11
View File
@@ -1,11 +0,0 @@
package monolith
import (
"testing"
"github.com/cortezaproject/corteza-server/pkg/app"
)
func TestConfigure(t *testing.T) {
var _ app.Runnable = &App{}
}
+1 -1
View File
@@ -12,8 +12,8 @@ import (
"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/options"
)
type (
-83
View File
@@ -1,83 +0,0 @@
package app
import (
"context"
"go.uber.org/zap"
"github.com/cortezaproject/corteza-server/pkg/actionlog"
)
// 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 {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Init)
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 {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Upgrade)
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 {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Activate)
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 {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Provision)
err = app.Provision(ctx)
if err != nil {
return
}
}
return
}
-29
View File
@@ -1,29 +0,0 @@
package options
import (
"time"
)
type (
GRPCServerOpt struct {
Network string `env:"GRPC_SERVER_NETWORK"`
Addr string `env:"GRPC_SERVER_ADDR"`
ClientMaxBackoffDelay time.Duration `env:"GRPC_CLIENT_BACKOFF_DELAY"`
ClientLog bool `env:"GRPC_CLIENT_LOG"`
}
)
func GRPCServer(pfix string) (o *GRPCServerOpt) {
o = &GRPCServerOpt{
Network: "tcp",
Addr: ":50051",
ClientMaxBackoffDelay: time.Minute,
ClientLog: false,
}
fill(o, pfix)
return
}
-122
View File
@@ -1,122 +0,0 @@
package options
import (
"os"
"reflect"
"strings"
"time"
"github.com/spf13/cast"
)
func fill(opt interface{}, pfix string) {
v := reflect.ValueOf(opt)
if v.Kind() != reflect.Ptr {
panic("expecting a pointer, not a value")
}
if v.IsNil() {
panic("nil pointer passed")
}
v = v.Elem()
length := v.NumField()
for i := 0; i < length; i++ {
f := v.Field(i)
t := v.Type().Field(i)
if tag := t.Tag.Get("env"); tag != "" {
if !f.CanSet() {
panic("unexpected pointer for field " + t.Name)
}
if f.Type() == reflect.TypeOf(time.Duration(1)) {
v.FieldByName(t.Name).SetInt(int64(EnvDuration(pfix, tag, time.Duration(f.Int()))))
continue
}
if f.Kind() == reflect.String {
v.FieldByName(t.Name).SetString(EnvString(pfix, tag, f.String()))
continue
}
if f.Kind() == reflect.Bool {
v.FieldByName(t.Name).SetBool(EnvBool(pfix, tag, f.Bool()))
continue
}
if f.Kind() == reflect.Int {
v.FieldByName(t.Name).SetInt(int64(EnvInt(pfix, tag, int(f.Int()))))
continue
}
if f.Kind() == reflect.Float32 {
v.FieldByName(t.Name).SetFloat(float64(EnvFloat32(pfix, tag, float32(f.Float()))))
continue
}
panic("unsupported type/kind for field " + t.Name)
}
}
}
func makeEnvKeys(pfix, name string) []string {
return []string{
strings.ToUpper(strings.Trim(pfix, "_") + "_" + name),
strings.ToUpper(name),
}
}
func EnvString(pfix, key string, def string) string {
for _, key = range makeEnvKeys(pfix, key) {
if val, has := os.LookupEnv(key); has {
return val
}
}
return def
}
func EnvBool(pfix, key string, def bool) bool {
for _, key = range makeEnvKeys(pfix, key) {
if val, has := os.LookupEnv(key); has {
if b, err := cast.ToBoolE(val); err == nil {
return b
}
}
}
return def
}
func EnvInt(pfix, key string, def int) int {
for _, key = range makeEnvKeys(pfix, key) {
if val, has := os.LookupEnv(key); has {
if i, err := cast.ToIntE(val); err == nil {
return i
}
}
}
return def
}
func EnvFloat32(pfix, key string, def float32) float32 {
for _, key = range makeEnvKeys(pfix, key) {
if val, has := os.LookupEnv(key); has {
if i, err := cast.ToFloat32E(val); err == nil {
return i
}
}
}
return def
}
func EnvDuration(pfix, key string, def time.Duration) time.Duration {
for _, key = range makeEnvKeys(pfix, key) {
if val, has := os.LookupEnv(key); has {
if d, err := cast.ToDurationE(val); err == nil {
return d
}
}
}
return def
}
-334
View File
@@ -1,334 +0,0 @@
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/actionlog"
"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(actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_API_REST))
wg.Done()
}(ctx)
}
if r.grpcServer != nil {
wg.Add(1)
go func(ctx context.Context) {
r.grpcServer.Serve(actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_API_GRPC))
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) {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Init)
if err = r.Initialize(ctx); err != nil {
return
}
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Upgrade)
if err = RunUpgrade(ctx, r.parts...); err != nil {
return
}
return
})
provisionCmd := cli.ProvisionCommand(func() (err error) {
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Activate)
if err = r.Activate(ctx); err != nil {
return
}
ctx = actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Provision)
if err = RunProvision(ctx, r.parts...); err != nil {
return
}
return
})
r.rootCmd.AddCommand(
serveCmd,
upgradeCmd,
provisionCmd,
cli.VersionCommand(),
)
// 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()
r.Run(actionlog.RequestOriginToContext(ctx, actionlog.RequestOrigin_APP_Run))
}
+18 -42
View File
@@ -9,30 +9,30 @@ import (
)
var (
rootCommand = cobra.Command{
rootCommand = &cobra.Command{
Use: "corteza-server",
Aliases: []string{"corteza", "server"},
TraverseChildren: true,
}
serveApiCommand = cobra.Command{
serveApiCommand = &cobra.Command{
Use: "serve-api",
Aliases: []string{"serve"},
Short: "Start HTTP server with REST API",
}
upgradeCommand = cobra.Command{
upgradeCommand = &cobra.Command{
Use: "upgrade",
Short: "Upgrade tasks",
}
provisionCommand = cobra.Command{
provisionCommand = &cobra.Command{
Use: "provision",
Short: "Provision tasks",
}
versionCommand = cobra.Command{
versionCommand = &cobra.Command{
Use: "version",
Short: "Version",
}
@@ -49,39 +49,24 @@ var (
//
// Callback is called when not executed with help subcommand
func RootCommand(ppRunEfn func() error) *cobra.Command {
// Make a copy
var (
cmd = rootCommand
silent, debug bool
)
cmd.PersistentPreRunE = func(cmd *cobra.Command, args []string) (err error) {
rootCommand.PersistentPreRunE = func(cmd *cobra.Command, args []string) (err error) {
if light[cmd.Name()] {
// Do not run hooks on parts on help command
return nil
}
if debug {
logger.DefaultLevel.SetLevel(zap.DebugLevel)
} else if silent {
logger.DefaultLevel.SetLevel(zap.FatalLevel)
if ppRunEfn != nil {
return ppRunEfn()
}
return ppRunEfn()
return nil
}
cmd.Flags().BoolVarP(&silent, "silent", "s", false, "No output")
cmd.Flags().BoolVarP(&debug, "debug", "d", false, "Debug")
return &cmd
return rootCommand
}
func ServeCommand(runEfn func() error) *cobra.Command {
// Make a copy
var cmd = serveApiCommand
cmd.RunE = func(cmd *cobra.Command, args []string) (err error) {
serveApiCommand.RunE = func(cmd *cobra.Command, args []string) (err error) {
if _, set := os.LookupEnv("LOG_LEVEL"); !set {
// If LOG_LEVEL is not explicitly set, let's
// set it to INFO so that it
@@ -91,38 +76,29 @@ func ServeCommand(runEfn func() error) *cobra.Command {
return runEfn()
}
return &cmd
return serveApiCommand
}
func UpgradeCommand(dbfn func() error) *cobra.Command {
// Make a copies
var dbCmd = upgradeCommand
dbCmd.RunE = func(cmd *cobra.Command, args []string) (err error) {
upgradeCommand.RunE = func(cmd *cobra.Command, args []string) (err error) {
return dbfn()
}
return &dbCmd
return upgradeCommand
}
func ProvisionCommand(cffn func() error) *cobra.Command {
// Make a copies
var cfCmd = provisionCommand
cfCmd.RunE = func(cmd *cobra.Command, args []string) (err error) {
provisionCommand.RunE = func(cmd *cobra.Command, args []string) (err error) {
return cffn()
}
return &cfCmd
return provisionCommand
}
func VersionCommand() *cobra.Command {
// Make a copies
var cfCmd = versionCommand
cfCmd.Run = func(cmd *cobra.Command, args []string) {
versionCommand.Run = func(cmd *cobra.Command, args []string) {
cmd.Println(version.Version)
}
return &cfCmd
return versionCommand
}
+1 -1
View File
@@ -11,8 +11,8 @@ import (
"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/pkg/options"
)
type (
+2 -2
View File
@@ -10,7 +10,7 @@ import (
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
"github.com/cortezaproject/corteza-server/pkg/app/options"
"github.com/cortezaproject/corteza-server/pkg/options"
)
var (
@@ -51,7 +51,7 @@ func Init() {
// Do we want to enable debug logger
// with a bit more dev-friendly output
debuggingLogger = options.EnvBool("", "LOG_DEBUG", false)
debuggingLogger = options.EnvBool("LOG_DEBUG", false)
)
if ll, has := os.LookupEnv("LOG_LEVEL"); has {
+1 -1
View File
@@ -8,7 +8,7 @@ import (
"go.uber.org/zap"
"github.com/cortezaproject/corteza-server/pkg/app/options"
"github.com/cortezaproject/corteza-server/pkg/options"
)
type Monitor struct {
@@ -13,7 +13,7 @@ func ActionLog() (o *ActionLogOpt) {
Debug: false,
}
fill(o, "")
fill(o)
return
}
@@ -51,7 +51,7 @@ func Corredor() (o *CorredorOpt) {
TlsCertPrivate: "private.key",
}
fill(o, "")
fill(o)
o.TlsCertCA = path.Join(o.TlsCertPath, o.TlsCertCA)
o.TlsCertPrivate = path.Join(o.TlsCertPath, o.TlsCertPrivate)
+8 -2
View File
@@ -1,6 +1,7 @@
package options
import (
"strings"
"time"
)
@@ -19,14 +20,19 @@ func DB(pfix string) (o *DBOpt) {
const maxTries = 100
o = &DBOpt{
DSN: "corteza:corteza@tcp(db:3306)/corteza?collation=utf8mb4_general_ci",
DSN: "mysql://corteza:corteza@tcp(db:3306)/corteza?collation=utf8mb4_general_ci",
Logger: false,
MaxTries: maxTries,
Delay: delay,
Timeout: maxTries * delay,
}
fill(o, pfix)
fill(o)
if !strings.Contains(o.DSN, "://") {
// Make sure DSN is compatible with new requirements
o.DSN = "mysql://" + o.DSN
}
return
}
+104
View File
@@ -0,0 +1,104 @@
package options
import (
"os"
"reflect"
"time"
"github.com/spf13/cast"
)
func fill(opt interface{}) {
v := reflect.ValueOf(opt)
if v.Kind() != reflect.Ptr {
panic("expecting a pointer, not a value")
}
if v.IsNil() {
panic("nil pointer passed")
}
v = v.Elem()
length := v.NumField()
for i := 0; i < length; i++ {
f := v.Field(i)
t := v.Type().Field(i)
if tag := t.Tag.Get("env"); tag != "" {
if !f.CanSet() {
panic("unexpected pointer for field " + t.Name)
}
if f.Type() == reflect.TypeOf(time.Duration(1)) {
v.FieldByName(t.Name).SetInt(int64(EnvDuration(tag, time.Duration(f.Int()))))
continue
}
if f.Kind() == reflect.String {
v.FieldByName(t.Name).SetString(EnvString(tag, f.String()))
continue
}
if f.Kind() == reflect.Bool {
v.FieldByName(t.Name).SetBool(EnvBool(tag, f.Bool()))
continue
}
if f.Kind() == reflect.Int {
v.FieldByName(t.Name).SetInt(int64(EnvInt(tag, int(f.Int()))))
continue
}
if f.Kind() == reflect.Float32 {
v.FieldByName(t.Name).SetFloat(float64(EnvFloat32(tag, float32(f.Float()))))
continue
}
panic("unsupported type/kind for field " + t.Name)
}
}
}
func EnvString(key string, def string) string {
if val, has := os.LookupEnv(key); has {
return val
}
return def
}
func EnvBool(key string, def bool) bool {
if val, has := os.LookupEnv(key); has {
if b, err := cast.ToBoolE(val); err == nil {
return b
}
}
return def
}
func EnvInt(key string, def int) int {
if val, has := os.LookupEnv(key); has {
if i, err := cast.ToIntE(val); err == nil {
return i
}
}
return def
}
func EnvFloat32(key string, def float32) float32 {
if val, has := os.LookupEnv(key); has {
if i, err := cast.ToFloat32E(val); err == nil {
return i
}
}
return def
}
func EnvDuration(key string, def time.Duration) time.Duration {
if val, has := os.LookupEnv(key); has {
if d, err := cast.ToDurationE(val); err == nil {
return d
}
}
return def
}
@@ -60,7 +60,7 @@ func HTTP(pfix string) (o *HTTPServerOpt) {
WebappList: "admin,auth,messaging,compose",
}
fill(o, pfix)
fill(o)
return
}
@@ -17,7 +17,7 @@ func HttpClient(pfix string) (o *HTTPClientOpt) {
HttpClientTimeout: 30 * time.Second,
}
fill(o, pfix)
fill(o)
return
}
@@ -18,7 +18,7 @@ func Auth() (o *AuthOpt) {
Expiry: time.Hour * 24 * 30,
}
fill(o, "")
fill(o)
// Setting JWT secret to random string to prevent security accidents...
//
@@ -14,7 +14,7 @@ func Monitor(pfix string) (o *MonitorOpt) {
o = &MonitorOpt{
Interval: 300 * time.Second,
}
fill(o, pfix)
fill(o)
return
}
+17
View File
@@ -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)
return
}
@@ -35,7 +35,7 @@ func PubSub(pfix string) (o *PubSubOpt) {
RedisPingPeriod: pingPeriod,
}
fill(o, pfix)
fill(o)
return
}
@@ -28,7 +28,7 @@ func Sentry(pfix string) (o *SentryOpt) {
Release: version.Version,
}
fill(o, pfix)
fill(o)
return
}
@@ -19,7 +19,7 @@ func SMTP(pfix string) (o *SMTPOpt) {
From: "",
}
fill(o, pfix)
fill(o)
return
}
@@ -26,7 +26,7 @@ func Storage(pfix string) (o *StorageOpt) {
MinioStrict: false,
}
fill(o, pfix)
fill(o)
return
}
@@ -2,16 +2,18 @@ package options
type (
UpgradeOpt struct {
Debug bool `env:"UPGRADE_DEBUG"`
Always bool `env:"UPGRADE_ALWAYS"`
}
)
func Upgrade(pfix string) (o *UpgradeOpt) {
o = &UpgradeOpt{
Debug: false,
Always: true,
}
fill(o, pfix)
fill(o)
return
}
@@ -26,7 +26,7 @@ func WaitFor(pfix string) (o *WaitForOpt) {
ServicesProbeInterval: time.Second * 5,
}
fill(o, pfix)
fill(o)
return
}
@@ -25,7 +25,7 @@ func Websocket(pfix string) (o *WebsocketOpt) {
PingPeriod: pingPeriod,
}
fill(o, pfix)
fill(o)
return
}
+1 -1
View File
@@ -3,8 +3,8 @@ package sentry
import (
"github.com/getsentry/sentry-go"
"github.com/cortezaproject/corteza-server/pkg/app/options"
"github.com/cortezaproject/corteza-server/pkg/logger"
"github.com/cortezaproject/corteza-server/pkg/options"
)
func Init(sentryOpt options.SentryOpt) error {
+1 -1
View File
@@ -2,7 +2,7 @@ package webapp
import (
"fmt"
"github.com/cortezaproject/corteza-server/pkg/app/options"
"github.com/cortezaproject/corteza-server/pkg/options"
"github.com/go-chi/chi"
"net/http"
"os"
-152
View File
@@ -1,152 +0,0 @@
package system
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/automation"
"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/pkg/scheduler"
"github.com/cortezaproject/corteza-server/system/auth/external"
"github.com/cortezaproject/corteza-server/system/commands"
migrate "github.com/cortezaproject/corteza-server/system/db"
systemGRPC "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"
"github.com/cortezaproject/corteza-server/system/service/event"
)
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
scheduler.Service().OnTick(
event.SystemOnInterval(),
event.SystemOnTimestamp(),
)
corredor.Service().SetUserFinder(service.DefaultUser)
corredor.Service().SetRoleFinder(service.DefaultRole)
return
}
func (app *App) Upgrade(ctx context.Context) (err error) {
db := factory.Database.MustGet().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{
ActionLog: app.Opts.ActionLog,
Storage: app.Opts.Storage,
})
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)
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, systemGRPC.NewUserService(
service.DefaultUser,
service.DefaultAuth,
auth.DefaultJwtHandler,
service.DefaultAccessControl,
))
proto.RegisterRolesServer(server, systemGRPC.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(),
commands.RBAC(),
// temp command, will be removed in 2020.6
automation.ScriptExporter(SERVICE),
)
}
-11
View File
@@ -1,11 +0,0 @@
package system
import (
"testing"
"github.com/cortezaproject/corteza-server/pkg/app"
)
func TestConfigure(t *testing.T) {
var _ app.Runnable = &App{}
}
+15 -14
View File
@@ -16,7 +16,7 @@ import (
)
// Will perform OpenID connect auto-configuration
func Auth() *cobra.Command {
func Auth(app serviceInitializer) *cobra.Command {
var (
enableDiscoveredProvider bool
skipValidationOnAutoDiscoveredProvider bool
@@ -28,14 +28,12 @@ func Auth() *cobra.Command {
}
autoDiscoverCmd := &cobra.Command{
Use: "auto-discovery [name] [url]",
Short: "Auto discovers new OIDC client",
Args: cobra.ExactArgs(2),
Use: "auto-discovery [name] [url]",
Short: "Auto discovers new OIDC client",
Args: cobra.ExactArgs(2),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
)
ctx := auth.SetSuperUserContext(cli.Context())
_, err := external.RegisterOidcProvider(
ctx,
args[0],
@@ -68,9 +66,10 @@ func Auth() *cobra.Command {
"Skip validation")
jwtCmd := &cobra.Command{
Use: "jwt [email-or-id]",
Short: "Generates new JWT for a user",
Args: cobra.MinimumNArgs(1),
Use: "jwt [email-or-id]",
Short: "Generates new JWT for a user",
Args: cobra.MinimumNArgs(1),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
@@ -111,10 +110,12 @@ func Auth() *cobra.Command {
}
testEmails := &cobra.Command{
Use: "test-notifications [recipient]",
Short: "Sends samples of all authentication notification to receipient",
Args: cobra.ExactArgs(1),
Use: "test-notifications [recipient]",
Short: "Sends samples of all authentication notification to receipient",
Args: cobra.ExactArgs(1),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
err error
+19
View File
@@ -0,0 +1,19 @@
package commands
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/cli"
"github.com/spf13/cobra"
)
type (
serviceInitializer interface {
InitServices(ctx context.Context) error
}
)
func commandPreRunInitService(app serviceInitializer) func(*cobra.Command, []string) error {
return func(_ *cobra.Command, _ []string) error {
return app.InitServices(cli.Context())
}
}
+6 -5
View File
@@ -50,24 +50,25 @@ type (
}
)
func RBAC() *cobra.Command {
func RBAC(app serviceInitializer) *cobra.Command {
cmd := &cobra.Command{
Use: "rbac",
Short: "RBAC tools",
Long: "Check and manipulates permissions",
}
cmd.AddCommand(rbacCheck())
cmd.AddCommand(rbacCheck(app))
//cmd.Flags().String("namespace", "", "Import into namespace (by ID or string)")
return cmd
}
func rbacCheck() *cobra.Command {
func rbacCheck(app serviceInitializer) *cobra.Command {
return &cobra.Command{
Use: "check",
Short: "Check applied permissions against given file (only supports compose permissions for now)",
Use: "check",
Short: "Check applied permissions against given file (only supports compose permissions for now)",
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
+5 -4
View File
@@ -13,16 +13,17 @@ import (
"github.com/cortezaproject/corteza-server/system/types"
)
func Roles() *cobra.Command {
func Roles(app serviceInitializer) *cobra.Command {
cmd := &cobra.Command{
Use: "roles",
Short: "Role management",
}
addUserCmd := &cobra.Command{
Use: "useradd [role-ID-or-name-or-handle] [user-ID-or-email]",
Short: "Add user to role",
Args: cobra.ExactArgs(2),
Use: "useradd [role-ID-or-name-or-handle] [user-ID-or-email]",
Short: "Add user to role",
Args: cobra.ExactArgs(2),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
// Create role and user repository.
var (
+7 -7
View File
@@ -2,6 +2,7 @@ package commands
import (
"encoding/json"
"github.com/cortezaproject/corteza-server/system/types"
"os"
"strings"
@@ -9,7 +10,6 @@ import (
"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"
)
@@ -47,7 +47,7 @@ func Settings() *cobra.Command {
},
}
list.Flags().String("prefix", "", "Filter settings by prefix")
list.Flags().String("prefix", "", "SettingsFilter settings by prefix")
get := &cobra.Command{
Use: "get [key to get, ...]",
@@ -79,7 +79,7 @@ func Settings() *cobra.Command {
value := args[1]
v := &settings.Value{
v := &types.SettingValue{
Name: args[0],
}
@@ -124,13 +124,13 @@ func Settings() *cobra.Command {
var (
decoder = json.NewDecoder(fh)
input = map[string]interface{}{}
vv settings.ValueSet
vv types.SettingValueSet
)
cli.HandleError(decoder.Decode(&input))
for k, v := range input {
val := &settings.Value{Name: k}
val := &types.SettingValue{Name: k}
cli.HandleError(val.SetValue(v))
vv = append(vv, val)
@@ -188,7 +188,7 @@ func Settings() *cobra.Command {
if vv, err := service.DefaultSettings.FindByPrefix(ctx); err != nil {
cli.HandleError(err)
} else {
_ = vv.Walk(func(v *settings.Value) error {
_ = vv.Walk(func(v *types.SettingValue) error {
names = append(names, v.Name)
return nil
})
@@ -203,7 +203,7 @@ func Settings() *cobra.Command {
},
}
del.Flags().String("prefix", "", "Filter settings by prefix")
del.Flags().String("prefix", "", "SettingsFilter settings by prefix")
cmd.AddCommand(
list,
+4 -3
View File
@@ -7,7 +7,7 @@ import (
)
// Will perform OpenID connect auto-configuration
func Sink() *cobra.Command {
func Sink(app serviceInitializer) *cobra.Command {
var (
expires string
srup = service.SinkRequestUrlParams{}
@@ -19,8 +19,9 @@ func Sink() *cobra.Command {
}
signatureCmd := &cobra.Command{
Use: "signature",
Short: "Creates signature for sink HTTP endpoint",
Use: "signature",
Short: "Creates signature for sink HTTP endpoint",
PreRunE: commandPreRunInitService(app),
RunE: func(cmd *cobra.Command, args []string) error {
if expires != "" {
// validate expiration date if set
+10 -6
View File
@@ -17,7 +17,7 @@ import (
"github.com/cortezaproject/corteza-server/system/types"
)
func Users() *cobra.Command {
func Users(app serviceInitializer) *cobra.Command {
var (
flagNoPassword bool
)
@@ -32,6 +32,8 @@ func Users() *cobra.Command {
listCmd := &cobra.Command{
Use: "list",
Short: "List users",
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
@@ -91,11 +93,12 @@ func Users() *cobra.Command {
Use: "add [email]",
Short: "Add new user",
Args: cobra.MinimumNArgs(1),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())
db = factory.Database.MustGet()
db = factory.Database.MustGet()
userRepo = repository.User(ctx, db)
authSvc = service.Auth(ctx)
@@ -146,9 +149,10 @@ func Users() *cobra.Command {
"Create user without password")
pwdCmd := &cobra.Command{
Use: "password [email]",
Short: "Change password for user",
Args: cobra.MinimumNArgs(1),
Use: "password [email]",
Short: "Change password for user",
Args: cobra.MinimumNArgs(1),
PreRunE: commandPreRunInitService(app),
Run: func(cmd *cobra.Command, args []string) {
var (
ctx = auth.SetSuperUserContext(cli.Context())