diff --git a/app/boot_levels.go b/app/boot_levels.go index f024894e5..4e55f3619 100644 --- a/app/boot_levels.go +++ b/app/boot_levels.go @@ -574,7 +574,7 @@ func (app *CortezaApp) InitServices(ctx context.Context) (err error) { sysService.DefaultReport.RegisterReporter("composeRecords", cmpService.DefaultRecord) // Initializing seeder - _ = seeder.Seeder(ctx, app.Store, seeder.Faker()) + _ = seeder.Seeder(ctx, app.Store, dal.Service(), seeder.Faker()) if err = app.plugins.Initialize(ctx, app.Log); err != nil { return fmt.Errorf("could not initialize plugins: %w", err) diff --git a/compose/rest/namespace.go b/compose/rest/namespace.go index f750b6fdc..d3662d352 100644 --- a/compose/rest/namespace.go +++ b/compose/rest/namespace.go @@ -18,6 +18,7 @@ import ( "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/api" "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/envoy" "github.com/cortezaproject/corteza-server/pkg/envoy/resource" envoyStore "github.com/cortezaproject/corteza-server/pkg/envoy/store" @@ -219,12 +220,12 @@ func (ctrl Namespace) Clone(ctx context.Context, r *request.NamespaceClone) (int decoder := func() (resource.InterfaceSet, error) { // get from store - return envoyStore.Decoder().Decode(ctx, service.DefaultStore, df) + return envoyStore.Decoder().Decode(ctx, service.DefaultStore, dal.Service(), df) } encoder := func(nn resource.InterfaceSet) error { // prepare for encoding - se := envoyStore.NewStoreEncoder(service.DefaultStore, &envoyStore.EncoderConfig{}) + se := envoyStore.NewStoreEncoder(service.DefaultStore, dal.Service(), &envoyStore.EncoderConfig{}) bld := envoy.NewBuilder(se) g, err := bld.Build(ctx, nn...) if err != nil { @@ -337,7 +338,7 @@ func (ctrl Namespace) ImportRun(ctx context.Context, r *request.NamespaceImportR } encoder = func(nn resource.InterfaceSet) error { - se := envoyStore.NewStoreEncoder(service.DefaultStore, &envoyStore.EncoderConfig{}) + se := envoyStore.NewStoreEncoder(service.DefaultStore, dal.Service(), &envoyStore.EncoderConfig{}) bld := envoy.NewBuilder(se) g, err := bld.Build(ctx, nn...) diff --git a/compose/rest/record.go b/compose/rest/record.go index 086ff8553..302cf20b5 100644 --- a/compose/rest/record.go +++ b/compose/rest/record.go @@ -16,6 +16,7 @@ import ( "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/api" "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/envoy" "github.com/cortezaproject/corteza-server/pkg/envoy/csv" ejson "github.com/cortezaproject/corteza-server/pkg/envoy/json" @@ -395,7 +396,7 @@ func (ctrl *Record) ImportRun(ctx context.Context, r *request.RecordImportRun) ( } return err } - se := estore.NewStoreEncoder(service.DefaultStore, cfg) + se := estore.NewStoreEncoder(service.DefaultStore, dal.Service(), cfg) bld := envoy.NewBuilder(se) g, err := bld.Build(ctx, ses.Resources...) if err != nil { @@ -464,7 +465,7 @@ func (ctrl *Record) Export(ctx context.Context, r *request.RecordExport) (interf } sd := estore.Decoder() - nn, err := sd.Decode(ctx, service.DefaultStore, f) + nn, err := sd.Decode(ctx, service.DefaultStore, dal.Service(), f) if err != nil { http.Error(w, fmt.Sprintf("failed to fetch records: %s", err.Error()), http.StatusBadRequest) } diff --git a/compose/service/attachment.go b/compose/service/attachment.go index 616ad3e29..deb663c8b 100644 --- a/compose/service/attachment.go +++ b/compose/service/attachment.go @@ -10,6 +10,7 @@ import ( "regexp" "strings" + "github.com/cortezaproject/corteza-server/compose/dalutils" "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/actionlog" "github.com/cortezaproject/corteza-server/pkg/auth" @@ -40,6 +41,7 @@ type ( objects objstore.Store ac attachmentAccessController store store.Storer + dal dalDater } attachmentAccessController interface { @@ -65,11 +67,12 @@ type ( } ) -func Attachment(store objstore.Store) *attachment { +func Attachment(store objstore.Store, dal dalDater) *attachment { return &attachment{ objects: store, ac: DefaultAccessControl, store: DefaultStore, + dal: dal, } } @@ -93,7 +96,7 @@ func (svc attachment) Find(ctx context.Context, filter types.AttachmentFilter) ( } if filter.RecordID > 0 { - aProps.namespace, aProps.module, aProps.record, err = loadRecordCombo(ctx, svc.store, filter.NamespaceID, filter.ModuleID, filter.RecordID) + aProps.namespace, aProps.module, aProps.record, err = loadRecordCombo(ctx, svc.store, svc.dal, filter.NamespaceID, filter.ModuleID, filter.RecordID) if err != nil { return err } else if svc.ac.CanReadRecord(ctx, aProps.record) { @@ -340,7 +343,7 @@ func (svc attachment) CreateRecordAttachment(ctx context.Context, namespaceID ui // To allow upload (attachment creation) user must have permissions to // alter that record - r, err = store.LookupComposeRecordByID(ctx, s, m, recordID) + r, err = dalutils.ComposeRecordsFind(ctx, svc.dal, m, recordID) if err != nil { return err } diff --git a/compose/service/record_datasource.go b/compose/service/record_datasource.go index 57037b721..6656cdf16 100644 --- a/compose/service/record_datasource.go +++ b/compose/service/record_datasource.go @@ -141,5 +141,6 @@ func (svc record) Datasource(ctx context.Context, ld *report.LoadStepDefinition) ld.Columns = cols } - return store.ComposeRecordDatasource(ctx, svc.store, mod, ld) + return nil, fmt.Errorf("@todo pending migration to DAL") + // return store.ComposeRecordDatasource(ctx, svc.store, mod, ld) } diff --git a/compose/service/service.go b/compose/service/service.go index 24ceec7f7..7b331c89c 100644 --- a/compose/service/service.go +++ b/compose/service/service.go @@ -178,7 +178,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config) DefaultPage = Page() DefaultChart = Chart() DefaultNotification = Notification(c.UserFinder) - DefaultAttachment = Attachment(DefaultObjectStore) + DefaultAttachment = Attachment(DefaultObjectStore, dal.Service()) RegisterIteratorProviders() diff --git a/pkg/envoy/store/compose.go b/pkg/envoy/store/compose.go index 1881ec2ee..346994e63 100644 --- a/pkg/envoy/store/compose.go +++ b/pkg/envoy/store/compose.go @@ -5,8 +5,11 @@ import ( "strconv" "strings" + "github.com/cortezaproject/corteza-server/compose/dalutils" "github.com/cortezaproject/corteza-server/compose/service" "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/dal" + "github.com/cortezaproject/corteza-server/pkg/dal/capabilities" "github.com/cortezaproject/corteza-server/pkg/envoy" "github.com/cortezaproject/corteza-server/pkg/envoy/resource" "github.com/cortezaproject/corteza-server/pkg/filter" @@ -29,19 +32,28 @@ type ( store.ComposeNamespaces store.ComposePages store.ComposeRecordValues - store.ComposeRecords + } + + dalService interface { + Create(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, vv ...dal.ValueGetter) error + Search(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, f filter.Filter) (dal.Iterator, error) + Delete(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv ...dal.ValueGetter) (err error) + Update(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv ...dal.ValueGetter) (err error) } composeDecoder struct { resourceID []uint64 namespaceID []uint64 + + dal dalService } ) -func newComposeDecoder() *composeDecoder { +func newComposeDecoder(dal dalService) *composeDecoder { return &composeDecoder{ resourceID: make([]uint64, 0, 200), namespaceID: make([]uint64, 0, 200), + dal: dal, } } @@ -252,7 +264,7 @@ func (d *composeDecoder) decodeComposeRecord(ctx context.Context, s store.Storer var err error for { - rr, fn, err = store.SearchComposeRecords(ctx, s, mod, types.RecordFilter(aux)) + rr, fn, err = dalutils.ComposeRecordsList(ctx, d.dal, mod, types.RecordFilter(aux)) if err != nil { return err } diff --git a/pkg/envoy/store/compose_record_marshal.go b/pkg/envoy/store/compose_record_marshal.go index d8dc57299..d9d073373 100644 --- a/pkg/envoy/store/compose_record_marshal.go +++ b/pkg/envoy/store/compose_record_marshal.go @@ -6,6 +6,7 @@ import ( "strconv" "time" + "github.com/cortezaproject/corteza-server/compose/dalutils" "github.com/cortezaproject/corteza-server/compose/service" "github.com/cortezaproject/corteza-server/compose/service/values" composeTypes "github.com/cortezaproject/corteza-server/compose/types" @@ -116,7 +117,7 @@ func (n *composeRecord) Prepare(ctx context.Context, pl *payload) (err error) { } // Preload all records - rr, _, err := store.SearchComposeRecords(ctx, pl.s, mod, composeTypes.RecordFilter{ + rr, _, err := dalutils.ComposeRecordsList(ctx, pl.dal, mod, composeTypes.RecordFilter{ ModuleID: mod.ID, NamespaceID: mod.NamespaceID, Paging: filter.Paging{ @@ -138,7 +139,7 @@ func (n *composeRecord) Prepare(ctx context.Context, pl *payload) (err error) { } // Preload own records - rr, _, err := store.SearchComposeRecords(ctx, pl.s, n.relMod, composeTypes.RecordFilter{ + rr, _, err := dalutils.ComposeRecordsList(ctx, pl.dal, n.relMod, composeTypes.RecordFilter{ ModuleID: n.relMod.ID, NamespaceID: n.relNS.ID, Paging: filter.Paging{ @@ -441,12 +442,12 @@ func (n *composeRecord) Encode(ctx context.Context, pl *payload) (err error) { // Create a new record if !exists { - err = store.CreateComposeRecord(ctx, pl.s, mod, rec) + err = dalutils.ComposeRecordCreate(ctx, pl.dal, mod, rec) return err } // Update existing - err = store.UpdateComposeRecord(ctx, pl.s, mod, rec) + err = dalutils.ComposeRecordUpdate(ctx, pl.dal, mod, rec) return err }() diff --git a/pkg/envoy/store/decoder.go b/pkg/envoy/store/decoder.go index db41e2e00..ac5df9673 100644 --- a/pkg/envoy/store/decoder.go +++ b/pkg/envoy/store/decoder.go @@ -107,7 +107,7 @@ func (aum auxMarshaller) MarshalEnvoy() ([]resource.Interface, error) { } // Decode decodes all of the things in the provided store -func (d *decoder) Decode(ctx context.Context, s store.Storer, f *DecodeFilter) ([]resource.Interface, error) { +func (d *decoder) Decode(ctx context.Context, s store.Storer, dal dalService, f *DecodeFilter) ([]resource.Interface, error) { mm := make(auxMarshaller, 0, 100) if d.ux == nil { @@ -134,7 +134,7 @@ func (d *decoder) Decode(ctx context.Context, s store.Storer, f *DecodeFilter) ( return mm, nil } - compose := newComposeDecoder() + compose := newComposeDecoder(dal) system := newSystemDecoder(d.ux) automation := newAutomationDecoder(d.ux) diff --git a/pkg/envoy/store/encoder.go b/pkg/envoy/store/encoder.go index 299f31a88..0e793fc6b 100644 --- a/pkg/envoy/store/encoder.go +++ b/pkg/envoy/store/encoder.go @@ -16,6 +16,7 @@ import ( type ( storeEncoder struct { s store.Storer + dal dalService cfg *EncoderConfig // Each resource should define its own state that is used when encoding the resource. @@ -61,6 +62,7 @@ type ( payload struct { s store.Storer + dal dalService state *envoy.ResourceState invokerID uint64 @@ -82,7 +84,7 @@ var ( // NewStoreEncoder initializes a fresh store encoder // // If no config is provided, it uses Skip as the default merge alg. -func NewStoreEncoder(s store.Storer, cfg *EncoderConfig) envoy.PrepareEncoder { +func NewStoreEncoder(s store.Storer, dal dalService, cfg *EncoderConfig) envoy.PrepareEncoder { if cfg == nil { cfg = &EncoderConfig{ OnExisting: resource.Skip, @@ -91,6 +93,7 @@ func NewStoreEncoder(s store.Storer, cfg *EncoderConfig) envoy.PrepareEncoder { return &storeEncoder{ s: s, + dal: dal, cfg: cfg, state: make(map[resource.Interface]resourceState), @@ -106,7 +109,7 @@ func (se *storeEncoder) Prepare(ctx context.Context, ee ...*envoy.ResourceState) return nil } - err = rs.Prepare(ctx, se.makePayload(ctx, se.s, ers)) + err = rs.Prepare(ctx, se.makePayload(ctx, se.s, se.dal, ers)) if err != nil { return err } @@ -192,7 +195,7 @@ func (se *storeEncoder) Encode(ctx context.Context, p envoy.Provider) error { err = ErrResourceStateUndefined } else { if state != nil { - err = state.Encode(ctx, se.makePayload(ctx, s, ers)) + err = state.Encode(ctx, se.makePayload(ctx, s, se.dal, ers)) } } @@ -203,9 +206,10 @@ func (se *storeEncoder) Encode(ctx context.Context, p envoy.Provider) error { }) } -func (se *storeEncoder) makePayload(ctx context.Context, s store.Storer, ers *envoy.ResourceState) *payload { +func (se *storeEncoder) makePayload(ctx context.Context, s store.Storer, dal dalService, ers *envoy.ResourceState) *payload { return &payload{ s: s, + dal: dal, state: ers, invokerID: auth.GetIdentityFromContext(ctx).Identity(), } diff --git a/pkg/seeder/commands/seeder.go b/pkg/seeder/commands/seeder.go index e66b18199..a2a4a5a9a 100644 --- a/pkg/seeder/commands/seeder.go +++ b/pkg/seeder/commands/seeder.go @@ -3,8 +3,10 @@ package commands import ( "context" "fmt" + "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/seeder" "github.com/spf13/cobra" @@ -39,7 +41,7 @@ func users(ctx context.Context, app serviceInitializer) (cmd *cobra.Command) { PreRunE: commandPreRunInitService(app), Run: func(cmd *cobra.Command, args []string) { var ( - seed = seeder.Seeder(ctx, seeder.DefaultStore, seeder.Faker()) + seed = seeder.Seeder(ctx, seeder.DefaultStore, dal.Service(), seeder.Faker()) err error ) @@ -90,7 +92,7 @@ func records(ctx context.Context, app serviceInitializer) (cmd *cobra.Command) { PreRunE: commandPreRunInitService(app), Run: func(cmd *cobra.Command, args []string) { var ( - seed = seeder.Seeder(ctx, seeder.DefaultStore, seeder.Faker()) + seed = seeder.Seeder(ctx, seeder.DefaultStore, dal.Service(), seeder.Faker()) err error ) @@ -145,7 +147,7 @@ func deleteAll(ctx context.Context, app serviceInitializer) (cmd *cobra.Command) PreRunE: commandPreRunInitService(app), Run: func(cmd *cobra.Command, args []string) { var ( - seed = seeder.Seeder(ctx, seeder.DefaultStore, seeder.Faker()) + seed = seeder.Seeder(ctx, seeder.DefaultStore, dal.Service(), seeder.Faker()) err error ) diff --git a/pkg/seeder/seeder.go b/pkg/seeder/seeder.go index 73376fddb..685590edb 100644 --- a/pkg/seeder/seeder.go +++ b/pkg/seeder/seeder.go @@ -3,10 +3,15 @@ package seeder import ( "context" "fmt" - cService "github.com/cortezaproject/corteza-server/compose/service" "time" + "github.com/cortezaproject/corteza-server/compose/dalutils" + cService "github.com/cortezaproject/corteza-server/compose/service" + cTypes "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/dal" + "github.com/cortezaproject/corteza-server/pkg/dal/capabilities" + "github.com/cortezaproject/corteza-server/pkg/filter" "github.com/cortezaproject/corteza-server/pkg/id" lTypes "github.com/cortezaproject/corteza-server/pkg/label/types" "github.com/cortezaproject/corteza-server/store" @@ -18,6 +23,7 @@ type ( ctx context.Context store storeService faker fakerService + dal dalService modSvc moduleService } @@ -50,10 +56,12 @@ type ( LookupComposeModuleByNamespaceIDHandle(ctx context.Context, namespaceID uint64, name string) (*cTypes.Module, error) LookupComposeModuleByID(ctx context.Context, id uint64) (*cTypes.Module, error) SearchComposeModuleFields(ctx context.Context, f cTypes.ModuleFieldFilter) (cTypes.ModuleFieldSet, cTypes.ModuleFieldFilter, error) + } - SearchComposeRecords(ctx context.Context, _mod *cTypes.Module, f cTypes.RecordFilter) (cTypes.RecordSet, cTypes.RecordFilter, error) - CreateComposeRecord(ctx context.Context, mod *cTypes.Module, rr ...*cTypes.Record) error - DeleteComposeRecord(ctx context.Context, m *cTypes.Module, rr ...*cTypes.Record) error + dalService interface { + Create(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, vv ...dal.ValueGetter) error + Search(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, f filter.Filter) (dal.Iterator, error) + Delete(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv ...dal.ValueGetter) (err error) } moduleService interface { @@ -71,13 +79,14 @@ const ( FakeDataLabel = "generated" ) -func Seeder(ctx context.Context, store store.Storer, faker fakerService) *seeder { +func Seeder(ctx context.Context, store store.Storer, dal dalService, faker fakerService) *seeder { DefaultStore = store return &seeder{ ctx: ctx, store: store, faker: faker, + dal: dal, modSvc: cService.DefaultModule, } } @@ -352,7 +361,8 @@ func (s seeder) CreateRecord(params RecordParams) (IDs []uint64, err error) { )) } - err = s.store.CreateComposeRecord(s.ctx, m, records...) + // @todo remove this .ctx ptrn + err = dalutils.ComposeRecordCreate(s.ctx, s.dal, m, records...) if err != nil { return } @@ -410,12 +420,12 @@ func (s seeder) DeleteAllRecord(mod *cTypes.Module) (err error) { FakeDataLabel: FakeDataLabel, }, } - records, _, err := s.store.SearchComposeRecords(s.ctx, mod, filter) + records, _, err := dalutils.ComposeRecordsList(s.ctx, s.dal, mod, filter) if err != nil { return } - err = s.store.DeleteComposeRecord(s.ctx, mod, records...) + err = dalutils.ComposeRecordDelete(s.ctx, s.dal, mod, records...) if err != nil { return } diff --git a/system/commands/exporter.go b/system/commands/exporter.go index aed5d2908..0ab5ab142 100644 --- a/system/commands/exporter.go +++ b/system/commands/exporter.go @@ -8,6 +8,7 @@ import ( "strings" "github.com/cortezaproject/corteza-server/pkg/auth" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/envoy/yaml" "github.com/spf13/cobra" @@ -40,7 +41,7 @@ func Export(ctx context.Context, storeInit func(ctx context.Context) (store.Stor f = f.FromResource(args...) sd := su.Decoder() - nn, err := sd.Decode(ctx, s, f) + nn, err := sd.Decode(ctx, s, dal.Service(), f) cli.HandleError(err) ye := yaml.NewYamlEncoder(&yaml.EncoderConfig{ diff --git a/system/commands/importer.go b/system/commands/importer.go index a6ef2da7b..a50f710aa 100644 --- a/system/commands/importer.go +++ b/system/commands/importer.go @@ -7,6 +7,7 @@ import ( "github.com/spf13/cobra" "github.com/cortezaproject/corteza-server/pkg/cli" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/envoy" "github.com/cortezaproject/corteza-server/pkg/envoy/csv" "github.com/cortezaproject/corteza-server/pkg/envoy/directory" @@ -76,7 +77,7 @@ func Import(ctx context.Context, storeInit func(ctx context.Context) (store.Stor opt.OnExisting = resource.MergeRight } - se := es.NewStoreEncoder(s, opt) + se := es.NewStoreEncoder(s, dal.Service(), opt) bld := envoy.NewBuilder(se) g, err := bld.Build(ctx, nn...) cli.HandleError(err) diff --git a/system/rest/user.go b/system/rest/user.go index 52d69dc93..6f5510909 100644 --- a/system/rest/user.go +++ b/system/rest/user.go @@ -13,6 +13,7 @@ import ( "github.com/cortezaproject/corteza-server/pkg/api" "github.com/cortezaproject/corteza-server/pkg/corredor" + "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/envoy" "github.com/cortezaproject/corteza-server/pkg/envoy/resource" envoyStore "github.com/cortezaproject/corteza-server/pkg/envoy/store" @@ -452,7 +453,7 @@ func (ctrl *User) Import(ctx context.Context, r *request.UserImport) (rsp interf } } - se := envoyStore.NewStoreEncoder(service.DefaultStore, &envoyStore.EncoderConfig{ + se := envoyStore.NewStoreEncoder(service.DefaultStore, dal.Service(), &envoyStore.EncoderConfig{ OnExisting: resource.Skip, })