Refactor core compose, system services with new DAL changes

* Define utility packages to work with DAL structs
* Cleanup code
This commit is contained in:
Tomaž Jerman
2022-06-14 12:08:16 +02:00
parent 8eb062293f
commit 033d2572dd
23 changed files with 1172 additions and 1279 deletions
+110
View File
@@ -0,0 +1,110 @@
package dalutils
import (
"context"
"github.com/cortezaproject/corteza-server/pkg/dal"
"github.com/cortezaproject/corteza-server/pkg/dal/capabilities"
"github.com/cortezaproject/corteza-server/store"
"github.com/cortezaproject/corteza-server/system/types"
)
type (
connectionCreator interface {
CreateConnection(ctx context.Context, connectionID uint64, cp dal.ConnectionParams, cm dal.ConnectionMeta, capabilities ...capabilities.Capability) (err error)
}
connectionDeleter interface {
DeleteConnection(ctx context.Context, connectionID uint64) (err error)
}
connectionUpdater interface {
UpdateConnection(ctx context.Context, connectionID uint64, cp dal.ConnectionParams, cm dal.ConnectionMeta, capabilities ...capabilities.Capability) (err error)
}
connectionDeleteCreator interface {
connectionDeleter
connectionCreator
}
)
func DalConnectionReload(ctx context.Context, s store.Storer, dc connectionDeleteCreator) (err error) {
// Get all available connections
cc, _, err := store.SearchDalConnections(ctx, s, types.DalConnectionFilter{
Type: types.DalConnectionResourceType,
})
if err != nil {
return
}
for _, c := range cc {
var cm dal.ConnectionMeta
cm, err = ConnectionMeta(ctx, c)
if err != nil {
return
}
if err = dc.CreateConnection(ctx, c.ID, c.Config.Connection, cm, c.ActiveCapabilities()...); err != nil {
return
}
}
return
}
func DalConnectionCreate(ctx context.Context, c connectionCreator, connections ...*types.DalConnection) (err error) {
var cm dal.ConnectionMeta
for _, connection := range connections {
cm, err = ConnectionMeta(ctx, connection)
if err != nil {
return
}
if err = c.CreateConnection(ctx, connection.ID, connection.Config.Connection, cm, connection.ActiveCapabilities()...); err != nil {
return err
}
}
return
}
func DalConnectionUpdate(ctx context.Context, u connectionUpdater, connections ...*types.DalConnection) (err error) {
var cm dal.ConnectionMeta
for _, connection := range connections {
cm, err = ConnectionMeta(ctx, connection)
if err != nil {
return
}
if err = u.UpdateConnection(ctx, connection.ID, connection.Config.Connection, cm, connection.ActiveCapabilities()...); err != nil {
return err
}
}
return
}
func DalConnectionDelete(ctx context.Context, d connectionDeleter, connections ...*types.DalConnection) (err error) {
for _, connection := range connections {
if err = d.DeleteConnection(ctx, connection.ID); err != nil {
return err
}
}
return
}
// // // // // // // // // // // // // // // // // // // // // // // // //
// Utils
func ConnectionMeta(ctx context.Context, c *types.DalConnection) (cm dal.ConnectionMeta, err error) {
// @todo we could probably utilize connection params more here
cm = dal.ConnectionMeta{
DefaultModelIdent: c.Config.DefaultModelIdent,
DefaultAttributeIdent: c.Config.DefaultAttributeIdent,
DefaultPartitionFormat: c.Config.DefaultPartitionFormat,
SensitivityLevel: c.SensitivityLevel,
Label: c.Handle,
}
return
}
+76
View File
@@ -0,0 +1,76 @@
package dalutils
import (
"context"
"sort"
"github.com/cortezaproject/corteza-server/pkg/dal"
"github.com/cortezaproject/corteza-server/pkg/filter"
"github.com/cortezaproject/corteza-server/store"
"github.com/cortezaproject/corteza-server/system/types"
)
type (
sensitivityLevelCreator interface {
CreateSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
}
sensitivityLevelUpdater interface {
UpdateSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
}
sensitivityLevelDeleter interface {
DeleteSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
}
sensitivityLevelReloader interface {
ReloadSensitivityLevels(levels ...dal.SensitivityLevel) (err error)
}
)
func DalSensitivityLevelReload(ctx context.Context, s store.Storer, r sensitivityLevelReloader) (err error) {
ll, _, err := store.SearchDalSensitivityLevels(ctx, s, types.DalSensitivityLevelFilter{Deleted: filter.StateExcluded})
if err != nil {
return
}
sort.Sort(ll)
levels := make(dal.SensitivityLevelSet, 0, len(ll))
for _, l := range ll {
levels = append(levels, systemToPkgType(l))
}
return r.ReloadSensitivityLevels(levels...)
}
func DalSensitivityLevelCreate(c sensitivityLevelCreator, levels ...*types.DalSensitivityLevel) (err error) {
return c.CreateSensitivityLevel(systemToPkgTypeSet(levels)...)
}
func DalSensitivityLevelUpdate(u sensitivityLevelUpdater, levels ...*types.DalSensitivityLevel) (err error) {
return u.UpdateSensitivityLevel(systemToPkgTypeSet(levels)...)
}
func DalSensitivityLevelDelete(d sensitivityLevelDeleter, levels ...*types.DalSensitivityLevel) (err error) {
return d.DeleteSensitivityLevel(systemToPkgTypeSet(levels)...)
}
// // // // // // // // // // // // // // // // // // // // // // // // //
// Utils
func systemToPkgType(level *types.DalSensitivityLevel) dal.SensitivityLevel {
return dal.SensitivityLevel{
ID: level.ID,
Handle: level.Handle,
Level: level.Level,
}
}
func systemToPkgTypeSet(levels types.DalSensitivityLevelSet) dal.SensitivityLevelSet {
out := make(dal.SensitivityLevelSet, 0, len(levels))
for _, l := range levels {
out = append(out, systemToPkgType(l))
}
return out
}
+47 -64
View File
@@ -12,6 +12,7 @@ import (
"github.com/cortezaproject/corteza-server/pkg/options"
"github.com/cortezaproject/corteza-server/store"
"github.com/cortezaproject/corteza-server/system/dalutils"
"github.com/cortezaproject/corteza-server/system/types"
)
@@ -35,9 +36,10 @@ type (
}
dalConnections interface {
AddConnection(ctx context.Context, connectionID uint64, cp dal.ConnectionParams, dft dal.ConnectionMeta, capabilities ...capabilities.Capability) (err error)
CreateConnection(ctx context.Context, connectionID uint64, cp dal.ConnectionParams, dft dal.ConnectionMeta, capabilities ...capabilities.Capability) (err error)
UpdateConnection(ctx context.Context, connectionID uint64, cp dal.ConnectionParams, dft dal.ConnectionMeta, capabilities ...capabilities.Capability) (err error)
RemoveConnection(ctx context.Context, connectionID uint64) (err error)
DeleteConnection(ctx context.Context, connectionID uint64) (err error)
SearchConnectionIssues(connectionID uint64) (err []error)
}
)
@@ -71,7 +73,7 @@ func (svc *dalConnection) FindByID(ctx context.Context, ID uint64) (q *types.Dal
return DalConnectionErrNotAllowedToRead(rProps)
}
svc.assurePrimaryConnection(q)
svc.proc(q)
return nil
}()
return q, svc.recordAction(ctx, rProps, DalConnectionActionLookup, err)
@@ -106,13 +108,11 @@ func (svc *dalConnection) Create(ctx context.Context, new *types.DalConnection)
q = new
var cm dal.ConnectionMeta
cm, err = svc.makeConnectionMeta(ctx, new)
if err != nil {
return
if err = dalutils.DalConnectionCreate(ctx, svc.dal, new); err != nil {
return err
}
return svc.dal.AddConnection(ctx, new.ID, new.Config.Connection, cm, new.ActiveCapabilities()...)
svc.proc(q)
return
}()
return q, svc.recordAction(ctx, qProps, DalConnectionActionCreate, err)
@@ -130,7 +130,7 @@ func (svc *dalConnection) Update(ctx context.Context, upd *types.DalConnection)
return DalConnectionErrNotFound(qProps)
}
svc.assurePrimaryConnection(old)
svc.proc(old)
if !svc.ac.CanUpdateDalConnection(ctx, old) {
return DalConnectionErrNotAllowedToUpdate(qProps)
@@ -175,15 +175,12 @@ func (svc *dalConnection) Update(ctx context.Context, upd *types.DalConnection)
}
q = upd
svc.assurePrimaryConnection(q)
var cm dal.ConnectionMeta
cm, err = svc.makeConnectionMeta(ctx, upd)
if err != nil {
return
defer svc.proc(q)
if old.HasIssues() {
return dalutils.DalConnectionCreate(ctx, svc.dal, upd)
}
return svc.dal.UpdateConnection(ctx, upd.ID, upd.Config.Connection, cm, upd.ActiveCapabilities()...)
return dalutils.DalConnectionUpdate(ctx, svc.dal, upd)
}()
return q, svc.recordAction(ctx, qProps, DalConnectionActionUpdate, err)
}
@@ -220,7 +217,7 @@ func (svc *dalConnection) DeleteByID(ctx context.Context, ID uint64) (err error)
return
}
return svc.dal.RemoveConnection(ctx, q.ID)
return dalutils.DalConnectionDelete(ctx, svc.dal, q)
}()
return svc.recordAction(ctx, qProps, DalConnectionActionDelete, err)
@@ -254,13 +251,8 @@ func (svc *dalConnection) UndeleteByID(ctx context.Context, ID uint64) (err erro
return
}
var cm dal.ConnectionMeta
cm, err = svc.makeConnectionMeta(ctx, q)
if err != nil {
return
}
return svc.dal.AddConnection(ctx, q.ID, q.Config.Connection, cm, q.ActiveCapabilities()...)
// We're creating it here since it was removed on delete
return dalutils.DalConnectionCreate(ctx, svc.dal, q)
}()
return svc.recordAction(ctx, qProps, DalConnectionActionDelete, err)
@@ -289,53 +281,44 @@ func (svc *dalConnection) Search(ctx context.Context, filter types.DalConnection
return err
}
svc.assurePrimaryConnection(r...)
svc.proc(r...)
return nil
}()
return r, f, svc.recordAction(ctx, aProps, DalConnectionActionSearch, err)
}
func (svc *dalConnection) ReloadConnections(ctx context.Context) (err error) {
// Get all available connections
cc, _, err := store.SearchDalConnections(ctx, svc.store, types.DalConnectionFilter{
Type: types.DalConnectionResourceType,
})
if err != nil {
return dalutils.DalConnectionReload(ctx, svc.store, svc.dal)
}
func (svc *dalConnection) proc(connections ...*types.DalConnection) {
for _, c := range connections {
svc.procPrimaryConnection(c)
svc.procDal(c)
svc.procLocale(c)
}
}
func (svc *dalConnection) procPrimaryConnection(c *types.DalConnection) {
if c.Type == types.DalPrimaryConnectionResourceType {
c.Config.Connection = dal.NewDSNConnection(svc.dbConf.DSN)
return
}
}
func (svc *dalConnection) procDal(c *types.DalConnection) {
ii := svc.dal.SearchConnectionIssues(c.ID)
if len(ii) == 0 {
c.Issues = nil
return
}
for _, c := range cc {
var cm dal.ConnectionMeta
cm, err = svc.makeConnectionMeta(ctx, c)
if err != nil {
return
}
if err = svc.dal.AddConnection(ctx, c.ID, c.Config.Connection, cm, c.ActiveCapabilities()...); err != nil {
return
}
}
return
}
func (svc *dalConnection) makeConnectionMeta(ctx context.Context, c *types.DalConnection) (cm dal.ConnectionMeta, err error) {
// @todo we could probably utilize connection params more here
cm = dal.ConnectionMeta{
DefaultModelIdent: c.Config.DefaultModelIdent,
DefaultAttributeIdent: c.Config.DefaultAttributeIdent,
DefaultPartitionFormat: c.Config.DefaultPartitionFormat,
SensitivityLevel: c.SensitivityLevel,
Label: c.Handle,
}
return
}
func (svc *dalConnection) assurePrimaryConnection(connections ...*types.DalConnection) {
for _, c := range connections {
if c.Type == types.DalPrimaryConnectionResourceType {
c.Config.Connection = dal.NewDSNConnection(svc.dbConf.DSN)
return
}
c.Issues = make([]string, len(ii))
for i, err := range ii {
c.Issues[i] = err.Error()
}
}
func (svc *dalConnection) procLocale(c *types.DalConnection) {
// @todo...
}
+47 -48
View File
@@ -8,9 +8,9 @@ import (
"github.com/cortezaproject/corteza-server/pkg/actionlog"
a "github.com/cortezaproject/corteza-server/pkg/auth"
"github.com/cortezaproject/corteza-server/pkg/dal"
"github.com/cortezaproject/corteza-server/pkg/filter"
"github.com/cortezaproject/corteza-server/store"
"github.com/cortezaproject/corteza-server/system/dalutils"
"github.com/cortezaproject/corteza-server/system/types"
)
@@ -27,7 +27,10 @@ type (
}
dalSensitivityLevels interface {
ReloadSensitivityLevels(raw dal.SensitivityLevelSet) (err error)
ReloadSensitivityLevels(levels ...dal.SensitivityLevel) (err error)
CreateSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
UpdateSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
DeleteSensitivityLevel(levels ...dal.SensitivityLevel) (err error)
}
)
@@ -72,28 +75,28 @@ func (svc *dalSensitivityLevel) Create(ctx context.Context, new *types.DalSensit
qProps = &dalSensitivityLevelActionProps{new: new}
)
err = func() (err error) {
err = store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) (err error) {
if !svc.ac.CanManageDalSensitivityLevel(ctx) {
return DalSensitivityLevelErrNotAllowedToManage(qProps)
}
new.CreatedAt = *now()
new.CreatedBy = a.GetIdentityFromContext(ctx).Identity()
ups, err := svc.prepare(ctx, svc.store, new)
ups, err := svc.prepare(ctx, s, new)
if err != nil {
return
}
new.ID = nextID()
err = store.UpsertDalSensitivityLevel(ctx, svc.store, ups...)
err = store.UpsertDalSensitivityLevel(ctx, s, ups...)
if err != nil {
return
}
q = new
return svc.ReloadSensitivityLevels(ctx, svc.store)
}()
return dalutils.DalSensitivityLevelCreate(svc.dal, new)
})
return q, svc.recordAction(ctx, qProps, DalSensitivityLevelActionCreate, err)
}
@@ -105,8 +108,8 @@ func (svc *dalSensitivityLevel) Update(ctx context.Context, upd *types.DalSensit
e error
)
err = func() (err error) {
if qq, e = store.LookupDalSensitivityLevelByID(ctx, svc.store, upd.ID); e != nil {
err = store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) (err error) {
if qq, e = store.LookupDalSensitivityLevelByID(ctx, s, upd.ID); e != nil {
return DalSensitivityLevelErrNotFound(qProps)
}
@@ -118,19 +121,19 @@ func (svc *dalSensitivityLevel) Update(ctx context.Context, upd *types.DalSensit
upd.CreatedAt = qq.CreatedAt
upd.UpdatedBy = a.GetIdentityFromContext(ctx).Identity()
ups, err := svc.prepare(ctx, svc.store, upd)
ups, err := svc.prepare(ctx, s, upd)
if err != nil {
return
}
err = store.UpsertDalSensitivityLevel(ctx, svc.store, ups...)
err = store.UpsertDalSensitivityLevel(ctx, s, ups...)
if err != nil {
return
}
q = upd
return svc.ReloadSensitivityLevels(ctx, svc.store)
}()
return dalutils.DalSensitivityLevelUpdate(svc.dal, upd)
})
return q, svc.recordAction(ctx, qProps, DalSensitivityLevelActionUpdate, err)
}
@@ -141,12 +144,12 @@ func (svc *dalSensitivityLevel) DeleteByID(ctx context.Context, ID uint64) (err
q *types.DalSensitivityLevel
)
err = func() (err error) {
err = store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) (err error) {
if ID == 0 {
return DalSensitivityLevelErrInvalidID()
}
if q, err = store.LookupDalSensitivityLevelByID(ctx, svc.store, ID); err != nil {
if q, err = store.LookupDalSensitivityLevelByID(ctx, s, ID); err != nil {
return
}
@@ -159,17 +162,36 @@ func (svc *dalSensitivityLevel) DeleteByID(ctx context.Context, ID uint64) (err
q.DeletedAt = now()
q.DeletedBy = a.GetIdentityFromContext(ctx).Identity()
ups, err := svc.prepare(ctx, svc.store, q)
ups, err := svc.prepare(ctx, s, q)
if err != nil {
return
}
err = store.UpsertDalSensitivityLevel(ctx, svc.store, ups...)
err = store.UpsertDalSensitivityLevel(ctx, s, ups...)
if err != nil {
return
}
return svc.ReloadSensitivityLevels(ctx, svc.store)
}()
var (
dd = make(types.DalSensitivityLevelSet, 0, len(ups)/2+1)
uu = make(types.DalSensitivityLevelSet, 0, len(ups)/2+1)
)
for _, l := range ups {
if l.DeletedAt != nil {
dd = append(dd, l)
} else {
uu = append(uu, l)
}
}
if err = dalutils.DalSensitivityLevelUpdate(svc.dal, uu...); err != nil {
return err
}
if err = dalutils.DalSensitivityLevelDelete(svc.dal, dd...); err != nil {
return err
}
return nil
})
return svc.recordAction(ctx, qProps, DalSensitivityLevelActionDelete, err)
}
@@ -180,12 +202,12 @@ func (svc *dalSensitivityLevel) UndeleteByID(ctx context.Context, ID uint64) (er
q *types.DalSensitivityLevel
)
err = func() (err error) {
err = store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) (err error) {
if ID == 0 {
return DalSensitivityLevelErrInvalidID()
}
if q, err = store.LookupDalSensitivityLevelByID(ctx, svc.store, ID); err != nil {
if q, err = store.LookupDalSensitivityLevelByID(ctx, s, ID); err != nil {
return
}
@@ -198,12 +220,12 @@ func (svc *dalSensitivityLevel) UndeleteByID(ctx context.Context, ID uint64) (er
q.DeletedAt = nil
q.UpdatedBy = a.GetIdentityFromContext(ctx).Identity()
if err = store.UpdateDalSensitivityLevel(ctx, svc.store, q); err != nil {
if err = store.UpdateDalSensitivityLevel(ctx, s, q); err != nil {
return
}
return svc.ReloadSensitivityLevels(ctx, svc.store)
}()
return dalutils.DalSensitivityLevelCreate(svc.dal, q)
})
return svc.recordAction(ctx, qProps, DalSensitivityLevelActionDelete, err)
}
@@ -238,30 +260,7 @@ func (svc *dalSensitivityLevel) Search(ctx context.Context, filter types.DalSens
}
func (svc *dalSensitivityLevel) ReloadSensitivityLevels(ctx context.Context, s store.Storer) (err error) {
ll, err := svc.getSensitivityLevels(ctx, s)
if err != nil {
return
}
return svc.dal.ReloadSensitivityLevels(ll)
}
func (svc *dalSensitivityLevel) getSensitivityLevels(ctx context.Context, s store.Storer) (out dal.SensitivityLevelSet, err error) {
ll, _, err := store.SearchDalSensitivityLevels(ctx, s, types.DalSensitivityLevelFilter{Deleted: filter.StateExcluded})
if err != nil {
return
}
sort.Sort(ll)
for _, l := range ll {
out = append(out, dal.SensitivityLevel{
ID: l.ID,
Handle: l.Handle,
})
}
return
return dalutils.DalSensitivityLevelReload(ctx, s, svc.dal)
}
func (svc *dalSensitivityLevel) prepare(ctx context.Context, s store.Storer, sl *types.DalSensitivityLevel) (_ types.DalSensitivityLevelSet, err error) {
+6
View File
@@ -24,6 +24,8 @@ type (
Ownership string `json:"ownership"`
SensitivityLevel uint64 `json:"sensitivityLevel,string"`
Issues []string `json:"issues,omitempty" db:"-"`
Config ConnectionConfig `json:"config"`
Capabilities ConnectionCapabilities `json:"capabilities"`
@@ -84,6 +86,10 @@ func (c DalConnection) ActiveCapabilities() capabilities.Set {
Union(c.Capabilities.Enabled)
}
func (c DalConnection) HasIssues() bool {
return len(c.Issues) > 0
}
func ParseConnectionConfig(ss []string) (m ConnectionConfig, err error) {
if len(ss) == 0 {
return