From 033d2572dd7d8214daceca2b90a70915fd4eba71 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Toma=C5=BE=20Jerman?= Date: Tue, 14 Jun 2022 11:01:02 +0200 Subject: [PATCH] Refactor core compose, system services with new DAL changes * Define utility packages to work with DAL structs * Cleanup code --- .../templates/gocode/store/interfaces.go.tpl | 1 - compose/dalutils/capabilities.go | 47 ++ .../module_dal.go => dalutils/modules.go} | 189 ++++--- compose/dalutils/records.go | 187 +++++++ compose/record.cue | 50 -- compose/service/dal_interfaces.go | 37 ++ compose/service/module.go | 78 ++- compose/service/record.go | 248 ++++----- compose/service/record_actions.gen.go | 276 ++++++++++ compose/service/record_actions.yaml | 30 ++ compose/service/record_dal.go | 123 ----- compose/types/module.go | 6 + store/adapters/rdbms/aux_types.gen.go | 67 --- store/adapters/rdbms/filters.gen.go | 23 - store/adapters/rdbms/queries.gen.go | 108 ---- store/adapters/rdbms/rdbms.gen.go | 485 ------------------ store/interfaces.gen.go | 95 ---- store/tests/all_test.go | 3 - system/dalutils/connection.go | 110 ++++ system/dalutils/sensitivity_level.go | 76 +++ system/service/dal_connection.go | 111 ++-- system/service/dal_sensitivity_level.go | 95 ++-- system/types/dal_connection.go | 6 + 23 files changed, 1172 insertions(+), 1279 deletions(-) create mode 100644 compose/dalutils/capabilities.go rename compose/{service/module_dal.go => dalutils/modules.go} (65%) create mode 100644 compose/dalutils/records.go create mode 100644 compose/service/dal_interfaces.go delete mode 100644 compose/service/record_dal.go create mode 100644 system/dalutils/connection.go create mode 100644 system/dalutils/sensitivity_level.go diff --git a/codegen/assets/templates/gocode/store/interfaces.go.tpl b/codegen/assets/templates/gocode/store/interfaces.go.tpl index 16079e9c2..578ef5830 100644 --- a/codegen/assets/templates/gocode/store/interfaces.go.tpl +++ b/codegen/assets/templates/gocode/store/interfaces.go.tpl @@ -10,7 +10,6 @@ import ( {{- end }} "github.com/cortezaproject/corteza-server/pkg/locale" "golang.org/x/text/language" - "github.com/cortezaproject/corteza-server/pkg/report" ) {{ define "extraArgs" -}} diff --git a/compose/dalutils/capabilities.go b/compose/dalutils/capabilities.go new file mode 100644 index 000000000..91a498c22 --- /dev/null +++ b/compose/dalutils/capabilities.go @@ -0,0 +1,47 @@ +package dalutils + +import ( + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/dal/capabilities" +) + +func recCreateCapabilities(m *types.Module) (out capabilities.Set) { + return capabilities.CreateCapabilities(m.ModelConfig.Capabilities...) +} + +func recUpdateCapabilities(m *types.Module) (out capabilities.Set) { + return capabilities.UpdateCapabilities(m.ModelConfig.Capabilities...) +} + +func recDeleteCapabilities(m *types.Module) (out capabilities.Set) { + return capabilities.DeleteCapabilities(m.ModelConfig.Capabilities...) +} + +func recFilterCapabilities(f types.RecordFilter) (out capabilities.Set) { + if f.PageCursor != nil { + out = append(out, capabilities.Paging) + } + + if f.IncPageNavigation { + out = append(out, capabilities.Paging) + } + + if f.IncTotal { + out = append(out, capabilities.Stats) + } + + if f.Sort != nil { + out = append(out, capabilities.Sorting) + } + + return +} + +func recSearchCapabilities(m *types.Module, f types.RecordFilter) (out capabilities.Set) { + return capabilities.SearchCapabilities(m.ModelConfig.Capabilities...). + Union(recFilterCapabilities(f)) +} + +func recLookupCapabilities(m *types.Module) (out capabilities.Set) { + return capabilities.LookupCapabilities(m.ModelConfig.Capabilities...) +} diff --git a/compose/service/module_dal.go b/compose/dalutils/modules.go similarity index 65% rename from compose/service/module_dal.go rename to compose/dalutils/modules.go index a0e3c3410..a0b432c4c 100644 --- a/compose/service/module_dal.go +++ b/compose/dalutils/modules.go @@ -1,4 +1,4 @@ -package service +package dalutils import ( "context" @@ -9,16 +9,38 @@ import ( "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/dal" "github.com/cortezaproject/corteza-server/pkg/handle" + "github.com/cortezaproject/corteza-server/store" systemTypes "github.com/cortezaproject/corteza-server/system/types" ) type ( - dalDDL interface { + identFormatter interface { ModelIdentFormatter(connectionID uint64) (f *dal.IdentFormatter, err error) + } + modelReloader interface { + identFormatter ReloadModel(ctx context.Context, models ...*dal.Model) (err error) - AddModel(ctx context.Context, models ...*dal.Model) (err error) - RemoveModel(ctx context.Context, models ...*dal.Model) (err error) + } + + modelCreator interface { + identFormatter + CreateModel(ctx context.Context, models ...*dal.Model) (err error) + } + + modelUpdater interface { + identFormatter + UpdateModel(ctx context.Context, old, new *dal.Model) (err error) + } + + attributeUpdater interface { + identFormatter + UpdateModelAttribute(ctx context.Context, model *dal.Model, old, new *dal.Attribute, trans ...dal.TransformationFunction) (err error) + } + + modelDeleter interface { + identFormatter + DeleteModel(ctx context.Context, models ...*dal.Model) (err error) } ) @@ -52,10 +74,7 @@ const ( colSysOwnedBy = "owned_by" ) -// ReloadDALModels reconstructs the DAL's data model based on the store.Storer -// -// Directly using store so we don't spam the action log -func (svc *module) ReloadDALModels(ctx context.Context) (err error) { +func ComposeModulesReload(ctx context.Context, s store.Storer, r modelReloader) (err error) { var ( namespaces types.NamespaceSet modules types.ModuleSet @@ -63,14 +82,14 @@ func (svc *module) ReloadDALModels(ctx context.Context) (err error) { models dal.ModelSet ) - namespaces, _, err = svc.store.SearchComposeNamespaces(ctx, types.NamespaceFilter{}) + namespaces, _, err = s.SearchComposeNamespaces(ctx, types.NamespaceFilter{}) if err != nil { return } var model *dal.Model for _, ns := range namespaces { - modules, _, err = svc.store.SearchComposeModules(ctx, types.ModuleFilter{ + modules, _, err = s.SearchComposeModules(ctx, types.ModuleFilter{ NamespaceID: ns.ID, }) if err != nil { @@ -78,7 +97,7 @@ func (svc *module) ReloadDALModels(ctx context.Context) (err error) { } for _, mod := range modules { - fields, _, err = svc.store.SearchComposeModuleFields(ctx, types.ModuleFieldFilter{ + fields, _, err = s.SearchComposeModuleFields(ctx, types.ModuleFieldFilter{ ModuleID: []uint64{mod.ID}, }) if err != nil { @@ -87,7 +106,7 @@ func (svc *module) ReloadDALModels(ctx context.Context) (err error) { mod.Fields = append(mod.Fields, fields...) - model, err = svc.moduleToModel(ctx, ns, mod) + model, err = moduleToModel(ctx, r, ns, mod) if err != nil { return } @@ -96,11 +115,91 @@ func (svc *module) ReloadDALModels(ctx context.Context) (err error) { } } - return svc.dal.ReloadModel(ctx, models...) + return r.ReloadModel(ctx, models...) } -func (svc *module) moduleToModel(ctx context.Context, ns *types.Namespace, mod *types.Module) (*dal.Model, error) { - formatter, tplParts, err := svc.prepareModelFormatter(ns, mod) +func ComposeModuleCreate(ctx context.Context, c modelCreator, ns *types.Namespace, mod *types.Module) (err error) { + model, err := moduleToModel(ctx, c, ns, mod) + if err != nil { + return + } + if err = c.CreateModel(ctx, model); err != nil { + return + } + + return +} + +func ComposeModuleUpdate(ctx context.Context, u modelUpdater, ns *types.Namespace, old, new *types.Module) (err error) { + oldModel, err := moduleToModel(ctx, u, ns, old) + if err != nil { + return + } + newModel, err := moduleToModel(ctx, u, ns, new) + if err != nil { + return + } + + if err = u.UpdateModel(ctx, oldModel, newModel); err != nil { + return + } + + return +} + +func ComposeModuleFieldsUpdate(ctx context.Context, u attributeUpdater, ns *types.Namespace, old, new *types.Module) (err error) { + oldModel, err := moduleToModel(ctx, u, ns, old) + if err != nil { + return + } + newModel, err := moduleToModel(ctx, u, ns, new) + if err != nil { + return + } + + diff := oldModel.Diff(newModel) + for _, d := range diff { + if err = u.UpdateModelAttribute(ctx, oldModel, d.Original, d.Asserted); err != nil { + return + } + } + + return +} + +func ComposeModuleDelete(ctx context.Context, d modelDeleter, ns *types.Namespace, mod *types.Module) (err error) { + model, err := moduleToModel(ctx, d, ns, mod) + if err != nil { + return + } + if err = d.DeleteModel(ctx, model); err != nil { + return + } + + return +} + +func ComposeModuleModelFormatter(f identFormatter, ns *types.Namespace, mod *types.Module) (formatter *dal.IdentFormatter, tplParts []string, err error) { + formatter, err = f.ModelIdentFormatter(mod.ModelConfig.ConnectionID) + if err != nil { + return + } + + modHandle, _ := handle.Cast(nil, mod.Handle, strconv.FormatUint(mod.ID, 10)) + nsHandle, _ := handle.Cast(nil, ns.Slug, strconv.FormatUint(ns.ID, 10)) + tplParts = []string{ + "module", modHandle, + "namespace", nsHandle, + } + + return +} + +// // // // // // // // // // // // // // // // // // // // // // // // // +// Utilities + +func moduleToModel(ctx context.Context, f identFormatter, ns *types.Namespace, mod *types.Module) (*dal.Model, error) { + formatter, tplParts, err := ComposeModuleModelFormatter(f, ns, mod) if err != nil { return nil, err } @@ -129,7 +228,7 @@ func (svc *module) moduleToModel(ctx context.Context, ns *types.Namespace, mod * // Handle user-defined fields for i, f := range mod.Fields { - out.Attributes[i], err = svc.moduleFieldToAttribute(getCodec, ns, mod, f) + out.Attributes[i], err = moduleFieldToAttribute(getCodec, ns, mod, f) if err != nil { return nil, err } @@ -138,16 +237,16 @@ func (svc *module) moduleToModel(ctx context.Context, ns *types.Namespace, mod * // Handle system fields; either default or user defined if !mod.ModelConfig.Partitioned { // When not partitioned the default system fields should be defined along side the `values` column - out.Attributes = append(out.Attributes, svc.moduleModelDefaultSysAttributes(getCodec)...) + out.Attributes = append(out.Attributes, moduleModelDefaultSysAttributes()...) } else { // When partitioned, we use store codec defined on the module - out.Attributes = append(out.Attributes, svc.moduleModelSysAttributes(mod, getCodec)...) + out.Attributes = append(out.Attributes, moduleModelSysAttributes(mod, getCodec)...) } return out, nil } -func (svc *module) moduleModelSysAttributes(mod *types.Module, getCodec func(f *types.ModuleField) dal.Codec) (out dal.AttributeSet) { +func moduleModelSysAttributes(mod *types.Module, getCodec func(f *types.ModuleField) dal.Codec) (out dal.AttributeSet) { sysEnc := mod.ModelConfig.SystemFieldEncoding if sysEnc.ID != nil { @@ -173,23 +272,23 @@ func (svc *module) moduleModelSysAttributes(mod *types.Module, getCodec func(f * } if sysEnc.UpdatedAt != nil { - out = append(out, dal.FullAttribute(sysUpdatedAt, &dal.TypeTimestamp{}, getCodec(&types.ModuleField{Name: sysUpdatedAt, EncodingStrategy: *sysEnc.UpdatedAt}))) + out = append(out, dal.FullAttribute(sysUpdatedAt, &dal.TypeTimestamp{Nullable: true}, getCodec(&types.ModuleField{Name: sysUpdatedAt, EncodingStrategy: *sysEnc.UpdatedAt}))) } if sysEnc.UpdatedBy != nil { - out = append(out, dal.FullAttribute(sysUpdatedBy, &dal.TypeID{}, getCodec(&types.ModuleField{Name: sysUpdatedBy, EncodingStrategy: *sysEnc.UpdatedBy}))) + out = append(out, dal.FullAttribute(sysUpdatedBy, &dal.TypeID{Nullable: true}, getCodec(&types.ModuleField{Name: sysUpdatedBy, EncodingStrategy: *sysEnc.UpdatedBy}))) } if sysEnc.DeletedAt != nil { - out = append(out, dal.FullAttribute(sysDeletedAt, &dal.TypeTimestamp{}, getCodec(&types.ModuleField{Name: sysDeletedAt, EncodingStrategy: *sysEnc.DeletedAt}))) + out = append(out, dal.FullAttribute(sysDeletedAt, &dal.TypeTimestamp{Nullable: true}, getCodec(&types.ModuleField{Name: sysDeletedAt, EncodingStrategy: *sysEnc.DeletedAt}))) } if sysEnc.DeletedBy != nil { - out = append(out, dal.FullAttribute(sysDeletedBy, &dal.TypeID{}, getCodec(&types.ModuleField{Name: sysDeletedBy, EncodingStrategy: *sysEnc.DeletedBy}))) + out = append(out, dal.FullAttribute(sysDeletedBy, &dal.TypeID{Nullable: true}, getCodec(&types.ModuleField{Name: sysDeletedBy, EncodingStrategy: *sysEnc.DeletedBy}))) } return } -func (svc *module) moduleModelDefaultSysAttributes(getCodec func(f *types.ModuleField) dal.Codec) dal.AttributeSet { +func moduleModelDefaultSysAttributes() dal.AttributeSet { return dal.AttributeSet{ dal.PrimaryAttribute(sysID, &dal.CodecAlias{Ident: colSysID}), @@ -201,15 +300,15 @@ func (svc *module) moduleModelDefaultSysAttributes(getCodec func(f *types.Module dal.FullAttribute(sysCreatedAt, &dal.TypeTimestamp{}, &dal.CodecAlias{Ident: colSysCreatedAt}), dal.FullAttribute(sysCreatedBy, &dal.TypeID{}, &dal.CodecAlias{Ident: colSysCreatedBy}), - dal.FullAttribute(sysUpdatedAt, &dal.TypeTimestamp{}, &dal.CodecAlias{Ident: colSysUpdatedAt}), - dal.FullAttribute(sysUpdatedBy, &dal.TypeID{}, &dal.CodecAlias{Ident: colSysUpdatedBy}), + dal.FullAttribute(sysUpdatedAt, &dal.TypeTimestamp{Nullable: true}, &dal.CodecAlias{Ident: colSysUpdatedAt}), + dal.FullAttribute(sysUpdatedBy, &dal.TypeID{Nullable: true}, &dal.CodecAlias{Ident: colSysUpdatedBy}), - dal.FullAttribute(sysDeletedAt, &dal.TypeTimestamp{}, &dal.CodecAlias{Ident: colSysDeletedAt}), - dal.FullAttribute(sysDeletedBy, &dal.TypeID{}, &dal.CodecAlias{Ident: colSysDeletedBy}), + dal.FullAttribute(sysDeletedAt, &dal.TypeTimestamp{Nullable: true}, &dal.CodecAlias{Ident: colSysDeletedAt}), + dal.FullAttribute(sysDeletedBy, &dal.TypeID{Nullable: true}, &dal.CodecAlias{Ident: colSysDeletedBy}), } } -func (svc *module) moduleFieldToAttribute(getCodec func(f *types.ModuleField) dal.Codec, ns *types.Namespace, mod *types.Module, f *types.ModuleField) (out *dal.Attribute, err error) { +func moduleFieldToAttribute(getCodec func(f *types.ModuleField) dal.Codec, ns *types.Namespace, mod *types.Module, f *types.ModuleField) (out *dal.Attribute, err error) { kind := f.Kind if kind == "" { kind = "String" @@ -326,35 +425,3 @@ func moduleFieldCodec(f *types.ModuleField, partitioned bool, formatter *dal.Ide return } - -// // // // // // // // // // // // // // // // // // // // // // // // // -// Utilities - -func (svc *module) addModuleToDAL(ctx context.Context, ns *types.Namespace, mod *types.Module) (err error) { - // Update DAL - model, err := svc.moduleToModel(ctx, ns, mod) - if err != nil { - return - } - if err = svc.dal.AddModel(ctx, model); err != nil { - return - } - - return -} - -func (svc *module) prepareModelFormatter(ns *types.Namespace, mod *types.Module) (formatter *dal.IdentFormatter, tplParts []string, err error) { - formatter, err = svc.dal.ModelIdentFormatter(mod.ModelConfig.ConnectionID) - if err != nil { - return - } - - modHandle, _ := handle.Cast(nil, mod.Handle, strconv.FormatUint(mod.ID, 10)) - nsHandle, _ := handle.Cast(nil, ns.Slug, strconv.FormatUint(ns.ID, 10)) - tplParts = []string{ - "module", modHandle, - "namespace", nsHandle, - } - - return -} diff --git a/compose/dalutils/records.go b/compose/dalutils/records.go new file mode 100644 index 000000000..56b533445 --- /dev/null +++ b/compose/dalutils/records.go @@ -0,0 +1,187 @@ +package dalutils + +import ( + "context" + "time" + + "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" +) + +type ( + creator interface { + Create(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, vv ...dal.ValueGetter) error + } + + updater interface { + Update(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, rr ...dal.ValueGetter) (err error) + } + + searcher interface { + Search(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, f filter.Filter) (dal.Iterator, error) + } + + lookuper interface { + Lookup(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, lookup dal.ValueGetter, dst dal.ValueSetter) (err error) + } + + deleter interface { + Delete(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv ...dal.ValueGetter) (err error) + } +) + +func ComposeRecordsList(ctx context.Context, s searcher, mod *types.Module, filter types.RecordFilter) (set types.RecordSet, outFilter types.RecordFilter, err error) { + iter, err := prepIterator(ctx, s, mod, filter) + if err != nil { + return + } + + set, outFilter, err = drainIterator(ctx, iter, mod, filter) + return +} + +func ComposeRecordsIterator(ctx context.Context, s searcher, mod *types.Module, filter types.RecordFilter) (iter dal.Iterator, outFilter types.RecordFilter, err error) { + iter, err = prepIterator(ctx, s, mod, filter) + if err != nil { + return + } + + outFilter = filter + outFilter.Paging = *filter.Paging.Clone() + + return +} + +func ComposeRecordsFind(ctx context.Context, l lookuper, mod *types.Module, recordID uint64) (out *types.Record, err error) { + out = prepareRecordTarget(mod) + + err = l.Lookup(ctx, mod.ModelFilter(), recLookupCapabilities(mod), dal.PKValues{"id": recordID}, out) + if err != nil { + return + } + + return +} + +func ComposeRecordCreate(ctx context.Context, c creator, mod *types.Module, records ...*types.Record) (err error) { + return c.Create(ctx, mod.ModelFilter(), recCreateCapabilities(mod), recToGetters(records...)...) +} + +func ComposeRecordUpdate(ctx context.Context, u updater, mod *types.Module, records ...*types.Record) (err error) { + return u.Update(ctx, mod.ModelFilter(), recUpdateCapabilities(mod), recToGetters(records...)...) +} + +func ComposeRecordSoftDelete(ctx context.Context, u updater, invoker uint64, mod *types.Module, records ...*types.Record) (err error) { + n := time.Now() + for _, r := range records { + r.DeletedAt = &n + r.DeletedBy = invoker + } + + return u.Update(ctx, mod.ModelFilter(), recUpdateCapabilities(mod), recToGetters(records...)...) +} + +func ComposeRecordDelete(ctx context.Context, d deleter, mod *types.Module, records ...*types.Record) (err error) { + return d.Delete(ctx, mod.ModelFilter(), recDeleteCapabilities(mod), recToGetters(records...)...) +} + +func WalkIterator(ctx context.Context, iter dal.Iterator, mod *types.Module, f func(r *types.Record) error) (err error) { + for iter.Next(ctx) { + r := prepareRecordTarget(mod) + if err = iter.Scan(r); err != nil { + return + } + + if err = f(r); err != nil { + return + } + } + + return iter.Err() +} + +// // // // // // // // // // // // // // // // // // // // // // // // // +// Utils + +func prepFilter(filter types.RecordFilter, mod *types.Module) (dalFilter filter.Filter) { + dalFilter = filter.ToFilter() + if mod.ModelConfig.Partitioned { + dalFilter = filter.ToConstraintedFilter(mod.ModelConfig.Constraints) + } + + return +} + +func prepIterator(ctx context.Context, dal searcher, mod *types.Module, filter types.RecordFilter) (iter dal.Iterator, err error) { + dalFilter := prepFilter(filter, mod) + + iter, err = dal.Search(ctx, mod.ModelFilter(), recSearchCapabilities(mod, filter), dalFilter) + return +} + +func drainIterator(ctx context.Context, iter dal.Iterator, mod *types.Module, filter types.RecordFilter) (set types.RecordSet, outFilter types.RecordFilter, err error) { + defer iter.Close() + + // Get the requested number of recrds + set = make(types.RecordSet, 0, filter.Limit) + i := 0 + err = WalkIterator(ctx, iter, mod, func(r *types.Record) error { + set = append(set, r) + i++ + return nil + }) + if err != nil { + return + } + + // Make out filter + outFilter = filter + pp := filter.Paging.Clone() + + if len(set) > 0 && filter.PrevPage != nil { + pp.PrevPage, err = iter.BackCursor(set[0]) + if err != nil { + return + } + } + + if len(set) > 0 { + pp.NextPage, err = iter.ForwardCursor(set[len(set)-1]) + if err != nil { + return + } + } + + outFilter.Paging = *pp + + return +} + +func prepareRecordTarget(module *types.Module) *types.Record { + // so we can avoid some code later involving (non)partitioned modules :seenoevil: + return &types.Record{ + ModuleID: module.ID, + NamespaceID: module.NamespaceID, + Values: make(types.RecordValueSet, 0, len(module.Fields)), + } +} + +func recToGetters(rr ...*types.Record) (out []dal.ValueGetter) { + out = make([]dal.ValueGetter, len(rr)) + + for i := range rr { + out[i] = rr[i] + } + + return +} + +func recToGetter(rr ...*types.Record) (out dal.ValueGetter) { + if len(rr) == 0 { + return + } + + return recToGetters(rr...)[0] +} diff --git a/compose/record.cue b/compose/record.cue index f5dd2426f..f10f66c1c 100644 --- a/compose/record.cue +++ b/compose/record.cue @@ -45,54 +45,4 @@ record: schema.#Resource & { "owner.manage": {} } } - - store: { - ident: "composeRecord" - - settings: { - rdbms: { - table: "compose_record" - } - } - - api: { - lookups: [ - { - fields: ["id"] - // export: false - description: """ - searches for compose record by ID - - It returns compose record even if deleted - """ - }, - ] - - functions: [ - { - expIdent: "ComposeRecordReport" - args: [ - { ident: "mod", goType: "*types.Module" }, - { ident: "metrics", goType: "string" }, - { ident: "dimensions", goType: "string" }, - { ident: "filters", goType: "string" } - ] - return: [ "[]map[string]interface{}" ] - }, { - expIdent: "ComposeRecordDatasource" - args: [ - { ident: "mod", goType: "*types.Module" }, - { ident: "ld", goType: "*report.LoadStepDefinition" } - ] - return: [ "report.Datasource" ] - }, { - expIdent: "PartialComposeRecordValueUpdate" - args: [ - { ident: "mod", goType: "*types.Module" }, - { ident: "values", goType: "*types.RecordValue", spread: true } - ] - } - ] - } - } } diff --git a/compose/service/dal_interfaces.go b/compose/service/dal_interfaces.go new file mode 100644 index 000000000..3c9b7b44e --- /dev/null +++ b/compose/service/dal_interfaces.go @@ -0,0 +1,37 @@ +package service + +import ( + "context" + + "github.com/cortezaproject/corteza-server/pkg/dal" + "github.com/cortezaproject/corteza-server/pkg/dal/capabilities" + "github.com/cortezaproject/corteza-server/pkg/filter" +) + +type ( + dalModeler interface { + ModelIdentFormatter(connectionID uint64) (f *dal.IdentFormatter, err error) + + ReloadModel(ctx context.Context, models ...*dal.Model) (err error) + CreateModel(ctx context.Context, models ...*dal.Model) (err error) + UpdateModel(ctx context.Context, old, new *dal.Model) (err error) + UpdateModelAttribute(ctx context.Context, model *dal.Model, old, new *dal.Attribute, trans ...dal.TransformationFunction) (err error) + DeleteModel(ctx context.Context, models ...*dal.Model) (err error) + + SearchModelIssues(connectionID, resourceID uint64) (out []error) + } + + dalDater interface { + Create(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, vv ...dal.ValueGetter) error + Update(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, rr ...dal.ValueGetter) (err error) + Search(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, f filter.Filter) (dal.Iterator, error) + Lookup(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, lookup dal.ValueGetter, dst dal.ValueSetter) (err error) + Delete(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv ...dal.ValueGetter) (err error) + Truncate(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set) (err error) + } + + dalService interface { + dalModeler + dalDater + } +) diff --git a/compose/service/module.go b/compose/service/module.go index d21b7912c..85034595e 100644 --- a/compose/service/module.go +++ b/compose/service/module.go @@ -7,13 +7,13 @@ import ( "sort" "strconv" + "github.com/cortezaproject/corteza-server/compose/dalutils" "github.com/cortezaproject/corteza-server/compose/service/event" "github.com/cortezaproject/corteza-server/compose/service/values" "github.com/cortezaproject/corteza-server/compose/types" "github.com/cortezaproject/corteza-server/pkg/actionlog" "github.com/cortezaproject/corteza-server/pkg/errors" "github.com/cortezaproject/corteza-server/pkg/eventbus" - "github.com/cortezaproject/corteza-server/pkg/filter" "github.com/cortezaproject/corteza-server/pkg/handle" "github.com/cortezaproject/corteza-server/pkg/label" "github.com/cortezaproject/corteza-server/pkg/locale" @@ -29,7 +29,7 @@ type ( store store.Storer locale ResourceTranslationsManagerService - dal dalDDL + dal dalService } moduleAccessController interface { @@ -82,7 +82,7 @@ var ( }) ) -func Module(ctx context.Context, dal dalDDL) (*module, error) { +func Module(ctx context.Context, dal dalService) (*module, error) { svc := &module{ ac: DefaultAccessControl, eventbus: eventbus.Service(), @@ -222,6 +222,11 @@ func (svc module) FindByAny(ctx context.Context, namespaceID uint64, identifier } func (svc module) proc(ctx context.Context, m *types.Module) { + svc.procLocale(ctx, m) + svc.procDal(m) +} + +func (svc module) procLocale(ctx context.Context, m *types.Module) { if svc.locale == nil || svc.locale.Locale() == nil { return } @@ -235,6 +240,23 @@ func (svc module) proc(ctx context.Context, m *types.Module) { }) } +func (svc module) procDal(m *types.Module) { + if svc.dal == nil { + return + } + + ii := svc.dal.SearchModelIssues(m.ModelConfig.ConnectionID, m.ID) + if len(ii) == 0 { + m.ModelConfig.Issues = nil + return + } + + m.ModelConfig.Issues = make([]string, len(ii)) + for i, err := range ii { + m.ModelConfig.Issues[i] = err.Error() + } +} + func (svc module) Create(ctx context.Context, new *types.Module) (*types.Module, error) { var ( ns *types.Namespace @@ -329,11 +351,13 @@ func (svc module) Create(ctx context.Context, new *types.Module) (*types.Module, return } - if err = svc.addModuleToDAL(ctx, ns, new); err != nil { + if err = dalutils.ComposeModuleCreate(ctx, svc.dal, ns, new); err != nil { return err } _ = svc.eventbus.WaitFor(ctx, event.ModuleAfterCreate(new, nil, ns)) + + svc.procDal(new) return nil }) @@ -352,6 +376,13 @@ func (svc module) UndeleteByID(ctx context.Context, namespaceID, moduleID uint64 return trim1st(svc.updater(ctx, namespaceID, moduleID, ModuleActionUndelete, svc.handleUndelete)) } +// ReloadDALModels reconstructs the DAL's data model based on the store.Storer +// +// Directly using store so we don't spam the action log +func (svc *module) ReloadDALModels(ctx context.Context) (err error) { + return dalutils.ComposeModulesReload(ctx, svc.store, svc.dal) +} + func (svc module) updater(ctx context.Context, namespaceID, moduleID uint64, action func(...*moduleActionProps) *moduleAction, fn moduleUpdateHandler) (*types.Module, error) { var ( changes moduleChanges @@ -373,6 +404,8 @@ func (svc module) updater(ctx context.Context, namespaceID, moduleID uint64, act } old = m.Clone() + // so we can get issues + svc.procDal(old) aProps.setNamespace(ns) aProps.setChanged(m) @@ -400,14 +433,15 @@ func (svc module) updater(ctx context.Context, namespaceID, moduleID uint64, act if changes&moduleFieldsChanged > 0 { var ( hasRecords bool - set types.RecordSet + // set types.RecordSet ) - if set, _, err = store.SearchComposeRecords(ctx, s, m, types.RecordFilter{Paging: filter.Paging{Limit: 1}}); err != nil { - return err - } + // @todo !! + // if set, _, err = dalutils.ComposeRecordsList(ctx, svc.dal, m, types.RecordFilter{Paging: filter.Paging{Limit: 1}}); err != nil { + // return err + // } - hasRecords = len(set) > 0 + // hasRecords = len(set) > 0 if err = updateModuleFields(ctx, s, m, old, hasRecords); err != nil { return err @@ -431,17 +465,28 @@ func (svc module) updater(ctx context.Context, namespaceID, moduleID uint64, act } if m.DeletedAt == nil { - err = svc.eventbus.WaitFor(ctx, event.ModuleAfterUpdate(m, old, ns)) + if err = svc.eventbus.WaitFor(ctx, event.ModuleAfterUpdate(m, old, ns)); err != nil { + return err + } + if err = dalutils.ComposeModuleUpdate(ctx, svc.dal, ns, old, m); err != nil { + return err + } + if err = dalutils.ComposeModuleFieldsUpdate(ctx, svc.dal, ns, old, m); err != nil { + return err + } } else { - err = svc.eventbus.WaitFor(ctx, event.ModuleAfterDelete(nil, old, ns)) + if err = svc.eventbus.WaitFor(ctx, event.ModuleAfterDelete(nil, old, ns)); err != nil { + return + } + if err = dalutils.ComposeModuleDelete(ctx, svc.dal, ns, m); err != nil { + return err + } } + svc.procDal(m) return err }) - // @todo improve this - err = svc.ReloadDALModels(ctx) - return m, svc.recordAction(ctx, aProps, action, err) } @@ -557,6 +602,11 @@ func (svc module) handleUpdate(ctx context.Context, upd *types.Module) moduleUpd res.ModelConfig = upd.ModelConfig } + if !reflect.DeepEqual(res.Privacy, upd.Privacy) { + changes |= moduleChanged + res.Privacy = upd.Privacy + } + // @todo make field-change detection more optimal if !reflect.DeepEqual(res.Fields, upd.Fields) { changes |= moduleFieldsChanged diff --git a/compose/service/record.go b/compose/service/record.go index 6887aeb80..fa9442f40 100644 --- a/compose/service/record.go +++ b/compose/service/record.go @@ -4,20 +4,21 @@ import ( "context" "encoding/json" "fmt" - "github.com/cortezaproject/corteza-server/pkg/locale" "regexp" "sort" "strconv" "time" + "github.com/cortezaproject/corteza-server/pkg/dal" + "github.com/cortezaproject/corteza-server/pkg/locale" + + "github.com/cortezaproject/corteza-server/compose/dalutils" "github.com/cortezaproject/corteza-server/compose/service/event" "github.com/cortezaproject/corteza-server/compose/service/values" "github.com/cortezaproject/corteza-server/compose/types" "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/dal" - "github.com/cortezaproject/corteza-server/pkg/dal/capabilities" "github.com/cortezaproject/corteza-server/pkg/envoy/resource" "github.com/cortezaproject/corteza-server/pkg/errors" "github.com/cortezaproject/corteza-server/pkg/eventbus" @@ -34,7 +35,7 @@ const ( type ( record struct { - dal dalDML + dal dalDater actionlog actionlog.Recorder @@ -156,7 +157,7 @@ type ( ErrorIndex map[string]int ) -func Record(dal dalDML) RecordService { +func Record(dal dalDater) RecordService { svc := &record{ actionlog: DefaultActionlog, ac: DefaultAccessControl, @@ -264,36 +265,16 @@ func (svc record) FindByID(ctx context.Context, namespaceID, moduleID, recordID return svc.lookup(ctx, namespaceID, moduleID, func(m *types.Module, props *recordActionProps) (*types.Record, error) { props.record.ID = recordID - out := svc.prepareRecordTarget(m) - return out, svc.dal.Lookup(ctx, m.ModelFilter(), capabilities.LookupCapabilities(m.ModelConfig.Capabilities...), dal.PKValues{"id": recordID}, out) + return dalutils.ComposeRecordsFind(ctx, svc.dal, m, recordID) }) } // Report generates report for a given module using metrics, dimensions and filter +// +// @todo remove; will be replaced with reporter endpoints func (svc record) Report(ctx context.Context, namespaceID, moduleID uint64, metrics, dimensions, filter string) (out interface{}, err error) { - var ( - ns *types.Namespace - m *types.Module - aProps = &recordActionProps{record: &types.Record{NamespaceID: namespaceID}} - ) - - err = func() error { - if ns, m, err = loadModuleWithNamespace(ctx, svc.store, namespaceID, moduleID); err != nil { - return err - } - - aProps.setNamespace(ns) - aProps.setModule(m) - - if !svc.ac.CanSearchRecordsOnModule(ctx, m) { - return RecordErrNotAllowedToSearch() - } - - out, err = store.ComposeRecordReport(ctx, svc.store, m, metrics, dimensions, filter) - return err - }() - - return out, svc.recordAction(ctx, aProps, RecordActionReport, err) + err = fmt.Errorf("record report endpoints pending removal") + return } func (svc record) Find(ctx context.Context, filter types.RecordFilter) (set types.RecordSet, f types.RecordFilter, err error) { @@ -331,19 +312,7 @@ func (svc record) Find(ctx context.Context, filter types.RecordFilter) (set type } } - dalFilter := filter.ToFilter() - if m.ModelConfig.Partitioned { - dalFilter = filter.ToConstraintedFilter(m.ModelConfig.Constraints) - } - - var iter dal.Iterator - iter, err = svc.dal.Search(ctx, m.ModelFilter(), svc.recSearchCapabilities(m, filter), dalFilter) - if err != nil { - return err - } - - set, f, err = svc.drainIterator(ctx, iter, filter, m) - if err != nil { + if set, f, err = dalutils.ComposeRecordsList(ctx, svc.dal, m, filter); err != nil { return err } @@ -542,15 +511,10 @@ func (svc record) create(ctx context.Context, new *types.Record) (rec *types.Rec return nil, RecordErrValueInput().Wrap(rve) } - err = svc.dal.Create(ctx, m.ModelFilter(), svc.recCreateCapabilities(m), svc.recToGetters(new)...) - if err != nil { + if err = dalutils.ComposeRecordCreate(ctx, svc.dal, m, new); err != nil { return } - if err != nil { - return nil, err - } - if err = label.Create(ctx, svc.store, new); err != nil { return } @@ -778,7 +742,7 @@ func (svc record) update(ctx context.Context, upd *types.Record) (rec *types.Rec return nil, RecordErrInvalidID() } - ns, m, old, err = loadRecordCombo(ctx, svc.store, upd.NamespaceID, upd.ModuleID, upd.ID) + ns, m, old, err = loadRecordCombo(ctx, svc.store, svc.dal, upd.NamespaceID, upd.ModuleID, upd.ID) if err != nil { return } @@ -838,7 +802,7 @@ func (svc record) update(ctx context.Context, upd *types.Record) (rec *types.Rec } } - return svc.dal.Update(ctx, m.ModelFilter(), svc.recUpdateCapabilities(m), svc.recToGetter(upd)) + return dalutils.ComposeRecordUpdate(ctx, svc.dal, m, upd) }) if err != nil { @@ -1022,7 +986,7 @@ func (svc record) delete(ctx context.Context, namespaceID, moduleID, recordID ui return nil, RecordErrInvalidID() } - ns, m, del, err = loadRecordCombo(ctx, svc.store, namespaceID, moduleID, recordID) + ns, m, del, err = loadRecordCombo(ctx, svc.store, svc.dal, namespaceID, moduleID, recordID) if err != nil { return nil, err } @@ -1041,11 +1005,7 @@ func (svc record) delete(ctx context.Context, namespaceID, moduleID, recordID ui } } - del.DeletedAt = now() - del.DeletedBy = invokerID - - err = svc.dal.Update(ctx, m.ModelFilter(), svc.recDeleteCapabilities(m), del) - if err != nil { + if err = dalutils.ComposeRecordSoftDelete(ctx, svc.dal, invokerID, m, del); err != nil { return nil, err } @@ -1129,15 +1089,13 @@ func (svc record) Organize(ctx context.Context, namespaceID, moduleID, recordID m *types.Module r *types.Record - recordValues = types.RecordValueSet{} - aProps = &recordActionProps{record: &types.Record{NamespaceID: namespaceID, ModuleID: moduleID, ID: recordID}} reorderingRecords bool ) err = func() error { - ns, m, r, err = loadRecordCombo(ctx, svc.store, namespaceID, moduleID, recordID) + ns, m, r, err = loadRecordCombo(ctx, svc.store, svc.dal, namespaceID, moduleID, recordID) if err != nil { return err } @@ -1153,77 +1111,75 @@ func (svc record) Organize(ctx context.Context, namespaceID, moduleID, recordID if posField != "" { reorderingRecords = true - if !regexp.MustCompile(`^[0-9]+$`).MatchString(position) { - return fmt.Errorf("expecting number for sorting position %q", posField) + // Position field checks + { + // Check field existence and permissions + // check if numeric -- we cannot reorder on any other field type + sf := m.Fields.FindByName(posField) + if sf == nil { + return RecordErrMissingPositionField() + } + + aProps.setPositionField(sf) + + if !sf.IsNumeric() { + return RecordErrInvalidPositionFieldKind() + } + + if sf.Multi { + return RecordErrInvalidPositionFieldConfigMultiValue() + } + + if !svc.ac.CanUpdateRecordValueOnModuleField(ctx, sf) { + return RecordErrNotAllowedToUpdate() + } } - // Check field existence and permissions - // check if numeric -- we cannot reorder on any other field type - - sf := m.Fields.FindByName(posField) - if sf == nil { - return fmt.Errorf("no such field %q", posField) - } - - if !sf.IsNumeric() { - return fmt.Errorf("cannot reorder on non numeric field %q", posField) - } - - if sf.Multi { - return fmt.Errorf("cannot reorder on multi-value field %q", posField) - } - - if !svc.ac.CanUpdateRecordValueOnModuleField(ctx, sf) { - return RecordErrNotAllowedToUpdate() + // Value checks + { + if !regexp.MustCompile(`^[0-9]+$`).MatchString(position) { + return RecordErrInvalidPositionValueType() + } } // Set new position - recordValues = recordValues.Set(&types.RecordValue{ - RecordID: recordID, - Name: posField, - Value: position, - }) + if err = r.SetValue(posField, 0, position); err != nil { + return err + } } if grpField != "" { - // Check field existence and permissions + // Group field checks + { + vf := m.Fields.FindByName(grpField) + if vf == nil { + return RecordErrMissingGroupField() + } - vf := m.Fields.FindByName(grpField) - if vf == nil { - return fmt.Errorf("no such field %q", grpField) - } + aProps.setGroupField(vf) - if vf.Multi { - return fmt.Errorf("cannot update multi-value field %q", posField) - } + if vf.Multi { + return RecordErrInvalidGroupFieldConfigMultiValue() + } - if !svc.ac.CanUpdateRecordValueOnModuleField(ctx, vf) { - return RecordErrNotAllowedToUpdate() + if !svc.ac.CanUpdateRecordValueOnModuleField(ctx, vf) { + return RecordErrNotAllowedToUpdate() + } } // Set new value - recordValues = recordValues.Set(&types.RecordValue{ - RecordID: recordID, - Name: grpField, - Value: group, - }) + if err = r.SetValue(grpField, 0, group); err != nil { + return err + } } return store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) error { - if len(recordValues) > 0 { - //svc.recordInfoUpdate(r) - //if err = store.UpdateComposeRecord(ctx, s, m, r); err != nil { - // return err - //} - } - - if err = store.PartialComposeRecordValueUpdate(ctx, s, m, recordValues...); err != nil { + if err = dalutils.ComposeRecordUpdate(ctx, svc.dal, m, r); err != nil { return err } if reorderingRecords { var ( - set types.RecordSet recordOrderPlace uint64 ) @@ -1246,26 +1202,40 @@ func (svc record) Organize(ctx context.Context, namespaceID, moduleID, recordID return err } - set, _, err = store.SearchComposeRecords(ctx, s, m, reorderFilter) + var iter dal.Iterator + if iter, _, err = dalutils.ComposeRecordsIterator(ctx, svc.dal, m, reorderFilter); err != nil { + return err + } + + const updateChunkSize = 300 + updateChunk := make(types.RecordSet, 0, updateChunkSize) + err = dalutils.WalkIterator(ctx, iter, m, func(r *types.Record) (err error) { + recordOrderPlace++ + if err = r.SetValue(posField, 0, strconv.FormatUint(recordOrderPlace, 10)); err != nil { + return + } + + // Update in chunks so we don't update one at a time and we don't keep + // too much data in memory. + updateChunk = append(updateChunk, r) + if len(updateChunk) >= updateChunkSize { + if err = dalutils.ComposeRecordUpdate(ctx, svc.dal, m, updateChunk...); err != nil { + return + } + updateChunk = make(types.RecordSet, 0, updateChunkSize) + } + + return + }) if err != nil { return err } - // Update value on each record - var vv = make([]*types.RecordValue, 0, len(set)) - _ = set.Walk(func(r *types.Record) error { - recordOrderPlace++ - vv = append(vv, &types.RecordValue{ - RecordID: r.ID, - Name: posField, - Value: strconv.FormatUint(recordOrderPlace, 10), - }) - - return nil - }) - - if err = store.PartialComposeRecordValueUpdate(ctx, s, m, vv...); err != nil { - return err + // Assure the last chunk is handled + if len(updateChunk) > 0 { + if err = dalutils.ComposeRecordUpdate(ctx, svc.dal, m, updateChunk...); err != nil { + return err + } } } @@ -1290,7 +1260,7 @@ func (svc record) Validate(ctx context.Context, rec *types.Record) error { // For backward compatibility (of controllers), it returns module+record func (svc record) TriggerScript(ctx context.Context, namespaceID, moduleID, recordID uint64, rvs types.RecordValueSet, script string) (*types.Module, *types.Record, error) { var ( - ns, m, r, err = loadRecordCombo(ctx, svc.store, namespaceID, moduleID, recordID) + ns, m, r, err = loadRecordCombo(ctx, svc.store, svc.dal, namespaceID, moduleID, recordID) ) if err != nil { @@ -1351,9 +1321,9 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu var ( invokerID = auth.GetIdentityFromContext(ctx).Identity() - ns *types.Namespace - m *types.Module - set types.RecordSet + ns *types.Namespace + m *types.Module + iter dal.Iterator aProps = &recordActionProps{} ) @@ -1364,13 +1334,12 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu return err } - // @todo might be good to split set into smaller chunks - set, f, err = store.SearchComposeRecords(ctx, svc.store, m, f) + iter, _, err = dalutils.ComposeRecordsIterator(ctx, svc.dal, m, f) if err != nil { return err } - for _, rec := range set { + err = dalutils.WalkIterator(ctx, iter, m, func(rec *types.Record) (err error) { switch action { case "clone": if !svc.ac.CanCreateRecordOnModule(ctx, m) { @@ -1416,7 +1385,7 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu } return store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) error { - return store.CreateComposeRecord(ctx, s, m, rec) + return dalutils.ComposeRecordCreate(ctx, svc.dal, m, rec) }) case "update": recordableAction = RecordActionIteratorUpdate @@ -1427,7 +1396,7 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu } return store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) error { - return store.UpdateComposeRecord(ctx, s, m, rec) + return dalutils.ComposeRecordUpdate(ctx, svc.dal, m, rec) }) case "delete": recordableAction = RecordActionIteratorDelete @@ -1435,7 +1404,7 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu return store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) error { rec.DeletedAt = now() rec.DeletedBy = invokerID - return store.UpdateComposeRecord(ctx, s, m, rec) + return dalutils.ComposeRecordSoftDelete(ctx, svc.dal, invokerID, m, rec) }) } @@ -1448,6 +1417,11 @@ func (svc record) Iterator(ctx context.Context, f types.RecordFilter, fn eventbu if err != nil { return err } + + return + }) + if err != nil { + return err } return nil @@ -1493,12 +1467,12 @@ func ComposeRecordFilterAC(ctx context.Context, ac recordValueAccessController, } // loadRecordCombo Loads namespace, module and record -func loadRecordCombo(ctx context.Context, s store.Storer, namespaceID, moduleID, recordID uint64) (ns *types.Namespace, m *types.Module, r *types.Record, err error) { +func loadRecordCombo(ctx context.Context, s store.Storer, dal dalDater, namespaceID, moduleID, recordID uint64) (ns *types.Namespace, m *types.Module, r *types.Record, err error) { if ns, m, err = loadModuleWithNamespace(ctx, s, namespaceID, moduleID); err != nil { return } - if r, err = store.LookupComposeRecordByID(ctx, s, m, recordID); err != nil { + if r, err = dalutils.ComposeRecordsFind(ctx, dal, m, recordID); err != nil { return } diff --git a/compose/service/record_actions.gen.go b/compose/service/record_actions.gen.go index 58d45b71b..0c8d0e357 100644 --- a/compose/service/record_actions.gen.go +++ b/compose/service/record_actions.gen.go @@ -28,6 +28,8 @@ type ( module *types.Module bulkOperation string field string + positionField *types.ModuleField + groupField *types.ModuleField value string valueErrors *types.RecordValueErrorSet } @@ -134,6 +136,28 @@ func (p *recordActionProps) setField(field string) *recordActionProps { return p } +// setPositionField updates recordActionProps's positionField +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *recordActionProps) setPositionField(positionField *types.ModuleField) *recordActionProps { + p.positionField = positionField + return p +} + +// setGroupField updates recordActionProps's groupField +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *recordActionProps) setGroupField(groupField *types.ModuleField) *recordActionProps { + p.groupField = groupField + return p +} + // setValue updates recordActionProps's value // // Allows method chaining @@ -198,6 +222,14 @@ func (p recordActionProps) Serialize() actionlog.Meta { } m.Set("bulkOperation", p.bulkOperation, true) m.Set("field", p.field, true) + if p.positionField != nil { + m.Set("positionField.name", p.positionField.Name, true) + m.Set("positionField.label", p.positionField.Label, true) + } + if p.groupField != nil { + m.Set("groupField.name", p.groupField.Name, true) + m.Set("groupField.label", p.groupField.Label, true) + } m.Set("value", p.value, true) if p.valueErrors != nil { m.Set("valueErrors.set", p.valueErrors.Set, true) @@ -324,6 +356,34 @@ func (p recordActionProps) Format(in string, err error) string { } pairs = append(pairs, "{{bulkOperation}}", fns(p.bulkOperation)) pairs = append(pairs, "{{field}}", fns(p.field)) + + if p.positionField != nil { + // replacement for "{{positionField}}" (in order how fields are defined) + pairs = append( + pairs, + "{{positionField}}", + fns( + p.positionField.Name, + p.positionField.Label, + ), + ) + pairs = append(pairs, "{{positionField.name}}", fns(p.positionField.Name)) + pairs = append(pairs, "{{positionField.label}}", fns(p.positionField.Label)) + } + + if p.groupField != nil { + // replacement for "{{groupField}}" (in order how fields are defined) + pairs = append( + pairs, + "{{groupField}}", + fns( + p.groupField.Name, + p.groupField.Label, + ), + ) + pairs = append(pairs, "{{groupField.name}}", fns(p.groupField.Name)) + pairs = append(pairs, "{{groupField.label}}", fns(p.groupField.Label)) + } pairs = append(pairs, "{{value}}", fns(p.value)) if p.valueErrors != nil { @@ -1330,6 +1390,222 @@ func RecordErrNotAllowedToChangeFieldValue(mm ...*recordActionProps) *errors.Err return e } +// RecordErrMissingPositionField returns "compose:record.missingPositionField" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrMissingPositionField(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("position module field not found", nil), + + errors.Meta("type", "missingPositionField"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "position module field not found"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.missingPositionField"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// RecordErrInvalidPositionFieldKind returns "compose:record.invalidPositionFieldKind" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrInvalidPositionFieldKind(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid position field {{positionField}} kind; kind must be 'Number'", nil), + + errors.Meta("type", "invalidPositionFieldKind"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "invalid position field {{positionField}} kind; kind must be 'Number'"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.invalidPositionFieldKind"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// RecordErrInvalidPositionFieldConfigMultiValue returns "compose:record.invalidPositionFieldConfigMultiValue" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrInvalidPositionFieldConfigMultiValue(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid position field {{positionField}} configuration; field must not be multi-value", nil), + + errors.Meta("type", "invalidPositionFieldConfigMultiValue"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "invalid position field {{positionField}} configuration; field must not be multi-value"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.invalidPositionFieldConfigMultiValue"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// RecordErrInvalidPositionValueType returns "compose:record.invalidPositionValueType" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrInvalidPositionValueType(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid position value data type; value must be numeric", nil), + + errors.Meta("type", "invalidPositionValueType"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "invalid position value data type; value must be numeric"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.invalidPositionValueType"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// RecordErrMissingGroupField returns "compose:record.missingGroupField" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrMissingGroupField(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("group module field not found", nil), + + errors.Meta("type", "missingGroupField"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "group module field not found"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.missingGroupField"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + +// RecordErrInvalidGroupFieldConfigMultiValue returns "compose:record.invalidGroupFieldConfigMultiValue" as *errors.Error +// +// +// This function is auto-generated. +// +func RecordErrInvalidGroupFieldConfigMultiValue(mm ...*recordActionProps) *errors.Error { + var p = &recordActionProps{} + if len(mm) > 0 { + p = mm[0] + } + + var e = errors.New( + errors.KindInternal, + + p.Format("invalid group field {{groupField}} configuration; field must not be multi-value", nil), + + errors.Meta("type", "invalidGroupFieldConfigMultiValue"), + errors.Meta("resource", "compose:record"), + + // action log entry; no formatting, it will be applied inside recordAction fn. + errors.Meta(recordLogMetaKey{}, "invalid group field {{groupField}} configuration; field must not be multi-value"), + errors.Meta(recordPropsMetaKey{}, p), + + // translation namespace & key + errors.Meta(locale.ErrorMetaNamespace{}, "compose"), + errors.Meta(locale.ErrorMetaKey{}, "record.errors.invalidGroupFieldConfigMultiValue"), + + errors.StackSkip(1), + ) + + if len(mm) > 0 { + } + + return e +} + // RecordErrImportSessionAlreadActive returns "compose:record.importSessionAlreadActive" as *errors.Error // // diff --git a/compose/service/record_actions.yaml b/compose/service/record_actions.yaml index 4bec41b84..3411806e0 100644 --- a/compose/service/record_actions.yaml +++ b/compose/service/record_actions.yaml @@ -30,6 +30,12 @@ props: fields: [ name, handle, ID, namespaceID ] - name: bulkOperation - name: field + - name: positionField + type: "*types.ModuleField" + fields: [ name, label ] + - name: groupField + type: "*types.ModuleField" + fields: [ name, label ] - name: value - name: valueErrors type: "*types.RecordValueErrorSet" @@ -156,7 +162,31 @@ errors: message: "not allowed to change value of field {{field}}" log: "failed to change value of field {{field}}; insufficient permissions" + # Organizer + - error: missingPositionField + message: "position module field not found" + log: "position module field not found" + - error: invalidPositionFieldKind + message: "invalid position field {{positionField}} kind; kind must be 'Number'" + log: "invalid position field {{positionField}} kind; kind must be 'Number'" + + - error: invalidPositionFieldConfigMultiValue + message: "invalid position field {{positionField}} configuration; field must not be multi-value" + log: "invalid position field {{positionField}} configuration; field must not be multi-value" + + - error: invalidPositionValueType + message: "invalid position value data type; value must be numeric" + log: "invalid position value data type; value must be numeric" + + - error: missingGroupField + message: "group module field not found" + log: "group module field not found" + - error: invalidGroupFieldConfigMultiValue + message: "invalid group field {{groupField}} configuration; field must not be multi-value" + log: "invalid group field {{groupField}} configuration; field must not be multi-value" + + # Importing - error: importSessionAlreadActive message: "import session already active" log: "failed to start import session" diff --git a/compose/service/record_dal.go b/compose/service/record_dal.go deleted file mode 100644 index 42d02e0f6..000000000 --- a/compose/service/record_dal.go +++ /dev/null @@ -1,123 +0,0 @@ -package service - -import ( - "context" - - "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" -) - -type ( - dalDML interface { - Create(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, vv ...dal.ValueGetter) error - Update(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, r dal.ValueGetter) (err error) - Search(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, f filter.Filter) (dal.Iterator, error) - Lookup(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, lookup dal.ValueGetter, dst dal.ValueSetter) (err error) - Delete(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set, pkv dal.ValueGetter) (err error) - Truncate(ctx context.Context, m dal.ModelFilter, capabilities capabilities.Set) (err error) - } -) - -func (svc *record) drainIterator(ctx context.Context, iter dal.Iterator, f types.RecordFilter, module *types.Module) (set types.RecordSet, outFilter types.RecordFilter, err error) { - set = make(types.RecordSet, 0, f.Limit) - - i := 0 - for iter.Next(ctx) { - auxr := svc.prepareRecordTarget(module) - if err = iter.Scan(auxr); err != nil { - return - } - - set = append(set, auxr) - - i++ - } - err = iter.Err() - - outFilter = f - pp := f.Paging.Clone() - - if len(set) > 0 && f.PrevPage != nil { - pp.PrevPage, err = iter.BackCursor(set[0]) - if err != nil { - return - } - } - - if len(set) > 0 { - pp.NextPage, err = iter.ForwardCursor(set[len(set)-1]) - if err != nil { - return - } - } - - outFilter.Paging = *pp - - return -} - -func (svc *record) prepareRecordTarget(module *types.Module) *types.Record { - // so we can avoid some code later involving (non)partitioned modules :seenoevil: - return &types.Record{ - ModuleID: module.ID, - NamespaceID: module.NamespaceID, - Values: make(types.RecordValueSet, 0, len(module.Fields)), - } -} - -func (svc *record) recToGetters(rr ...*types.Record) (out []dal.ValueGetter) { - out = make([]dal.ValueGetter, len(rr)) - - for i := range rr { - out[i] = rr[i] - } - - return -} - -func (svc *record) recToGetter(rr ...*types.Record) (out dal.ValueGetter) { - if len(rr) == 0 { - return - } - - return svc.recToGetters(rr...)[0] -} - -func (svc *record) recCreateCapabilities(m *types.Module) (out capabilities.Set) { - return capabilities.CreateCapabilities(m.ModelConfig.Capabilities...) -} - -func (svc *record) recUpdateCapabilities(m *types.Module) (out capabilities.Set) { - return capabilities.UpdateCapabilities(m.ModelConfig.Capabilities...) -} - -func (svc *record) recDeleteCapabilities(m *types.Module) (out capabilities.Set) { - return capabilities.DeleteCapabilities(m.ModelConfig.Capabilities...) -} - -func (svc *record) recFilterCapabilities(f types.RecordFilter) (out capabilities.Set) { - if f.PageCursor != nil { - out = append(out, capabilities.Paging) - } - - if f.IncPageNavigation { - out = append(out, capabilities.Paging) - } - - if f.IncTotal { - out = append(out, capabilities.Stats) - } - - if f.Sort != nil { - out = append(out, capabilities.Sorting) - } - - return -} - -func (svc *record) recSearchCapabilities(m *types.Module, f types.RecordFilter) (out capabilities.Set) { - return capabilities.SearchCapabilities(m.ModelConfig.Capabilities...). - Union(svc.recFilterCapabilities(f)) -} diff --git a/compose/types/module.go b/compose/types/module.go index c33a430cd..0e8f6a0be 100644 --- a/compose/types/module.go +++ b/compose/types/module.go @@ -19,6 +19,8 @@ type ( ConnectionID uint64 `json:"connectionID,string"` Capabilities capabilities.Set `json:"capabilities"` + Issues []string `json:"issues,omitempty"` + Constraints map[string][]any `json:"constraints"` Partitioned bool `json:"partitioned"` @@ -108,6 +110,10 @@ func (m Module) Clone() *Module { return c } +func (m Module) HasIssues() bool { + return len(m.ModelConfig.Issues) > 0 +} + // We won't worry about fields at this point func (m *Module) decodeTranslations(tt locale.ResourceTranslationIndex) { return diff --git a/store/adapters/rdbms/aux_types.gen.go b/store/adapters/rdbms/aux_types.gen.go index 86246b8e0..564581df6 100644 --- a/store/adapters/rdbms/aux_types.gen.go +++ b/store/adapters/rdbms/aux_types.gen.go @@ -308,20 +308,6 @@ type ( DeletedAt *time.Time `db:"deleted_at"` } - // auxComposeRecord is an auxiliary structure used for transporting to/from RDBMS store - auxComposeRecord struct { - ID uint64 `db:"id"` - ModuleID uint64 `db:"module_id"` - NamespaceID uint64 `db:"namespace_id"` - OwnedBy uint64 `db:"owned_by"` - CreatedAt time.Time `db:"created_at"` - UpdatedAt *time.Time `db:"updated_at"` - DeletedAt *time.Time `db:"deleted_at"` - CreatedBy uint64 `db:"created_by"` - UpdatedBy uint64 `db:"updated_by"` - DeletedBy uint64 `db:"deleted_by"` - } - // auxComposeRecordValue is an auxiliary structure used for transporting to/from RDBMS store auxComposeRecordValue struct { RecordID uint64 `db:"record_id"` @@ -1700,59 +1686,6 @@ func (aux *auxComposePage) scan(row scanner) error { ) } -// encodes ComposeRecord to auxComposeRecord -// -// This function is auto-generated -func (aux *auxComposeRecord) encode(res *composeType.Record) (_ error) { - aux.ID = res.ID - aux.ModuleID = res.ModuleID - aux.NamespaceID = res.NamespaceID - aux.OwnedBy = res.OwnedBy - aux.CreatedAt = res.CreatedAt - aux.UpdatedAt = res.UpdatedAt - aux.DeletedAt = res.DeletedAt - aux.CreatedBy = res.CreatedBy - aux.UpdatedBy = res.UpdatedBy - aux.DeletedBy = res.DeletedBy - return -} - -// decodes ComposeRecord from auxComposeRecord -// -// This function is auto-generated -func (aux auxComposeRecord) decode() (res *composeType.Record, _ error) { - res = new(composeType.Record) - res.ID = aux.ID - res.ModuleID = aux.ModuleID - res.NamespaceID = aux.NamespaceID - res.OwnedBy = aux.OwnedBy - res.CreatedAt = aux.CreatedAt - res.UpdatedAt = aux.UpdatedAt - res.DeletedAt = aux.DeletedAt - res.CreatedBy = aux.CreatedBy - res.UpdatedBy = aux.UpdatedBy - res.DeletedBy = aux.DeletedBy - return -} - -// scans row and fills auxComposeRecord fields -// -// This function is auto-generated -func (aux *auxComposeRecord) scan(row scanner) error { - return row.Scan( - &aux.ID, - &aux.ModuleID, - &aux.NamespaceID, - &aux.OwnedBy, - &aux.CreatedAt, - &aux.UpdatedAt, - &aux.DeletedAt, - &aux.CreatedBy, - &aux.UpdatedBy, - &aux.DeletedBy, - ) -} - // encodes ComposeRecordValue to auxComposeRecordValue // // This function is auto-generated diff --git a/store/adapters/rdbms/filters.gen.go b/store/adapters/rdbms/filters.gen.go index 15e2d8ab9..028823df5 100644 --- a/store/adapters/rdbms/filters.gen.go +++ b/store/adapters/rdbms/filters.gen.go @@ -83,9 +83,6 @@ type ( // optional composePage filter function called after the generated function ComposePage func(*Store, composeType.PageFilter) ([]goqu.Expression, composeType.PageFilter, error) - // optional composeRecord filter function called after the generated function - ComposeRecord func(*Store, composeType.RecordFilter) ([]goqu.Expression, composeType.RecordFilter, error) - // optional composeRecordValue filter function called after the generated function ComposeRecordValue func(*Store, composeType.RecordValueFilter) ([]goqu.Expression, composeType.RecordValueFilter, error) @@ -661,26 +658,6 @@ func ComposePageFilter(f composeType.PageFilter) (ee []goqu.Expression, _ compos return ee, f, err } -// ComposeRecordFilter returns logical expressions -// -// This function is called from Store.QueryComposeRecords() and can be extended -// by setting Store.Filters.ComposeRecord. Extension is called after all expressions -// are generated and can choose to ignore or alter them. -// -// This function is auto-generated -func ComposeRecordFilter(f composeType.RecordFilter) (ee []goqu.Expression, _ composeType.RecordFilter, err error) { - - if expr := stateNilComparison("deleted_at", f.Deleted); expr != nil { - ee = append(ee, expr) - } - - if len(f.LabeledIDs) > 0 { - ee = append(ee, goqu.I("id").In(f.LabeledIDs)) - } - - return ee, f, err -} - // ComposeRecordValueFilter returns logical expressions // // This function is called from Store.QueryComposeRecordValues() and can be extended diff --git a/store/adapters/rdbms/queries.gen.go b/store/adapters/rdbms/queries.gen.go index 5bf49eb31..c5efc768c 100644 --- a/store/adapters/rdbms/queries.gen.go +++ b/store/adapters/rdbms/queries.gen.go @@ -2095,114 +2095,6 @@ var ( } } - // composeRecordTable represents composeRecords store table - // - // This value is auto-generated - composeRecordTable = goqu.T("compose_record") - - // composeRecordSelectQuery assembles select query for fetching composeRecords - // - // This function is auto-generated - composeRecordSelectQuery = func(d goqu.DialectWrapper) *goqu.SelectDataset { - return d.Select( - "id", - "module_id", - "rel_namespace", - "owned_by", - "created_at", - "updated_at", - "deleted_at", - "created_by", - "updated_by", - "deleted_by", - ).From(composeRecordTable) - } - - // composeRecordInsertQuery assembles query inserting composeRecords - // - // This function is auto-generated - composeRecordInsertQuery = func(d goqu.DialectWrapper, res *composeType.Record) *goqu.InsertDataset { - return d.Insert(composeRecordTable). - Rows(goqu.Record{ - "id": res.ID, - "module_id": res.ModuleID, - "rel_namespace": res.NamespaceID, - "owned_by": res.OwnedBy, - "created_at": res.CreatedAt, - "updated_at": res.UpdatedAt, - "deleted_at": res.DeletedAt, - "created_by": res.CreatedBy, - "updated_by": res.UpdatedBy, - "deleted_by": res.DeletedBy, - }) - } - - // composeRecordUpsertQuery assembles (insert+on-conflict) query for replacing composeRecords - // - // This function is auto-generated - composeRecordUpsertQuery = func(d goqu.DialectWrapper, res *composeType.Record) *goqu.InsertDataset { - var target = `,id` - - return composeRecordInsertQuery(d, res). - OnConflict( - goqu.DoUpdate(target[1:], - goqu.Record{ - "module_id": res.ModuleID, - "rel_namespace": res.NamespaceID, - "owned_by": res.OwnedBy, - "created_at": res.CreatedAt, - "updated_at": res.UpdatedAt, - "deleted_at": res.DeletedAt, - "created_by": res.CreatedBy, - "updated_by": res.UpdatedBy, - "deleted_by": res.DeletedBy, - }, - ), - ) - } - - // composeRecordUpdateQuery assembles query for updating composeRecords - // - // This function is auto-generated - composeRecordUpdateQuery = func(d goqu.DialectWrapper, res *composeType.Record) *goqu.UpdateDataset { - return d.Update(composeRecordTable). - Set(goqu.Record{ - "module_id": res.ModuleID, - "rel_namespace": res.NamespaceID, - "owned_by": res.OwnedBy, - "created_at": res.CreatedAt, - "updated_at": res.UpdatedAt, - "deleted_at": res.DeletedAt, - "created_by": res.CreatedBy, - "updated_by": res.UpdatedBy, - "deleted_by": res.DeletedBy, - }). - Where(composeRecordPrimaryKeys(res)) - } - - // composeRecordDeleteQuery assembles delete query for removing composeRecords - // - // This function is auto-generated - composeRecordDeleteQuery = func(d goqu.DialectWrapper, ee ...goqu.Expression) *goqu.DeleteDataset { - return d.Delete(composeRecordTable).Where(ee...) - } - - // composeRecordDeleteQuery assembles delete query for removing composeRecords - // - // This function is auto-generated - composeRecordTruncateQuery = func(d goqu.DialectWrapper) *goqu.TruncateDataset { - return d.Truncate(composeRecordTable) - } - - // composeRecordPrimaryKeys assembles set of conditions for all primary keys - // - // This function is auto-generated - composeRecordPrimaryKeys = func(res *composeType.Record) goqu.Ex { - return goqu.Ex{ - "id": res.ID, - } - } - // composeRecordValueTable represents composeRecordValues store table // // This value is auto-generated diff --git a/store/adapters/rdbms/rdbms.gen.go b/store/adapters/rdbms/rdbms.gen.go index a4fd95691..a70f81ec1 100644 --- a/store/adapters/rdbms/rdbms.gen.go +++ b/store/adapters/rdbms/rdbms.gen.go @@ -47,7 +47,6 @@ var ( _ store.ComposeModuleFields = &Store{} _ store.ComposeNamespaces = &Store{} _ store.ComposePages = &Store{} - _ store.ComposeRecords = &Store{} _ store.ComposeRecordValues = &Store{} _ store.Credentials = &Store{} _ store.DalConnections = &Store{} @@ -8520,490 +8519,6 @@ func (s *Store) checkComposePageConstraints(ctx context.Context, res *composeTyp return nil } -// CreateComposeRecord creates one or more rows in composeRecord collection -// -// This function is auto-generated -func (s *Store) CreateComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) (err error) { - for i := range rr { - if err = s.checkComposeRecordConstraints(ctx, rr[i]); err != nil { - return - } - - if err = s.Exec(ctx, composeRecordInsertQuery(s.Dialect, rr[i])); err != nil { - return - } - } - - return -} - -// UpdateComposeRecord updates one or more existing entries in composeRecord collection -// -// This function is auto-generated -func (s *Store) UpdateComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) (err error) { - for i := range rr { - if err = s.checkComposeRecordConstraints(ctx, rr[i]); err != nil { - return - } - - if err = s.Exec(ctx, composeRecordUpdateQuery(s.Dialect, rr[i])); err != nil { - return - } - } - - return -} - -// UpsertComposeRecord updates one or more existing entries in composeRecord collection -// -// This function is auto-generated -func (s *Store) UpsertComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) (err error) { - for i := range rr { - if err = s.checkComposeRecordConstraints(ctx, rr[i]); err != nil { - return - } - - if err = s.Exec(ctx, composeRecordUpsertQuery(s.Dialect, rr[i])); err != nil { - return - } - } - - return -} - -// DeleteComposeRecord Deletes one or more entries from composeRecord collection -// -// This function is auto-generated -func (s *Store) DeleteComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) (err error) { - for i := range rr { - if err = s.Exec(ctx, composeRecordDeleteQuery(s.Dialect, composeRecordPrimaryKeys(rr[i]))); err != nil { - return - } - } - - return nil -} - -// DeleteComposeRecordByID deletes single entry from composeRecord collection -// -// This function is auto-generated -func (s *Store) DeleteComposeRecordByID(ctx context.Context, mod *composeType.Module, id uint64) error { - return s.Exec(ctx, composeRecordDeleteQuery(s.Dialect, goqu.Ex{ - "id": id, - })) -} - -// TruncateComposeRecords Deletes all rows from the composeRecord collection -func (s Store) TruncateComposeRecords(ctx context.Context, mod *composeType.Module) error { - return s.Exec(ctx, composeRecordTruncateQuery(s.Dialect)) -} - -// SearchComposeRecords returns (filtered) set of ComposeRecords -// -// This function is auto-generated -func (s *Store) SearchComposeRecords(ctx context.Context, mod *composeType.Module, f composeType.RecordFilter) (set composeType.RecordSet, _ composeType.RecordFilter, err error) { - - // Cleanup unwanted cursor values (only relevant is f.PageCursor, next&prev are reset and returned) - f.PrevPage, f.NextPage = nil, nil - - if f.PageCursor != nil { - // Page cursor exists; we need to validate it against used sort - // To cover the case when paging cursor is set but sorting is empty, we collect the sorting instructions - // from the cursor. - // This (extracted sorting info) is then returned as part of response - if f.Sort, err = f.PageCursor.Sort(f.Sort); err != nil { - return - } - } - - // Make sure results are always sorted at least by primary keys - if f.Sort.Get("id") == nil { - f.Sort = append(f.Sort, &filter.SortExpr{ - Column: "id", - Descending: f.Sort.LastDescending(), - }) - } - - // Cloned sorting instructions for the actual sorting - // Original are passed to the etchFullPageOfComposeRecords fn used for cursor creation; - // direction information it MUST keep the initial - sort := f.Sort.Clone() - - // When cursor for a previous page is used it's marked as reversed - // This tells us to flip the descending flag on all used sort keys - if f.PageCursor != nil && f.PageCursor.ROrder { - sort.Reverse() - } - - set, f.PrevPage, f.NextPage, err = s.fetchFullPageOfComposeRecords(ctx, f, sort) - - f.PageCursor = nil - if err != nil { - return nil, f, err - } - - return set, f, nil -} - -// fetchFullPageOfComposeRecords collects all requested results. -// -// Function applies: -// - cursor conditions (where ...) -// - limit -// -// Main responsibility of this function is to perform additional sequential queries in case when not enough results -// are collected due to failed check on a specific row (by check fn). -// -// Function then moves cursor to the last item fetched -// -// This function is auto-generated -func (s *Store) fetchFullPageOfComposeRecords( - ctx context.Context, - filter composeType.RecordFilter, - sort filter.SortExprSet, -) (set []*composeType.Record, prev, next *filter.PagingCursor, err error) { - var ( - aux []*composeType.Record - - // When cursor for a previous page is used it's marked as reversed - // This tells us to flip the descending flag on all used sort keys - reversedOrder = filter.PageCursor != nil && filter.PageCursor.ROrder - - // Copy no. of required items to limit - // Limit will change when doing subsequent queries to fill - // the set with all required items - limit = filter.Limit - - reqItems = filter.Limit - - // cursor to prev. page is only calculated when cursor is used - hasPrev = filter.PageCursor != nil - - // next cursor is calculated when there are more pages to come - hasNext bool - - tryFilter composeType.RecordFilter - ) - - set = make([]*composeType.Record, 0, DefaultSliceCapacity) - - for try := 0; try < MaxRefetches; try++ { - // Copy filter & apply custom sorting that might be affected by cursor - tryFilter = filter - tryFilter.Sort = sort - - if limit > 0 { - // fetching + 1 to peak ahead if there are more items - // we can fetch (next-page cursor) - tryFilter.Limit = limit + 1 - } - - if aux, hasNext, err = s.QueryComposeRecords(ctx, tryFilter); err != nil { - return nil, nil, nil, err - } - - if len(aux) == 0 { - // nothing fetched - break - } - - // append fetched items - set = append(set, aux...) - - if reqItems == 0 || !hasNext { - // no max requested items specified, break out - break - } - - collected := uint(len(set)) - - if reqItems > collected { - // not enough items fetched, try again with adjusted limit - limit = reqItems - collected - - if limit < MinEnsureFetchLimit { - // In case limit is set very low and we've missed records in the first fetch, - // make sure next fetch limit is a bit higher - limit = MinEnsureFetchLimit - } - - // Update cursor so that it points to the last item fetched - tryFilter.PageCursor = s.collectComposeRecordCursorValues(set[collected-1], filter.Sort...) - - // Copy reverse flag from sorting - tryFilter.PageCursor.LThen = filter.Sort.Reversed() - continue - } - - if reqItems < collected { - set = set[:reqItems] - } - - break - } - - collected := len(set) - - if collected == 0 { - return nil, nil, nil, nil - } - - if reversedOrder { - // Fetched set needs to be reversed because we've forced a descending order to get the previous page - for i, j := 0, collected-1; i < j; i, j = i+1, j-1 { - set[i], set[j] = set[j], set[i] - } - - // when in reverse-order rules on what cursor to return change - hasPrev, hasNext = hasNext, hasPrev - } - - if hasPrev { - prev = s.collectComposeRecordCursorValues(set[0], filter.Sort...) - prev.ROrder = true - prev.LThen = !filter.Sort.Reversed() - } - - if hasNext { - next = s.collectComposeRecordCursorValues(set[collected-1], filter.Sort...) - next.LThen = filter.Sort.Reversed() - } - - return set, prev, next, nil -} - -// QueryComposeRecords queries the database, converts and checks each row and returns collected set -// -// With generics, we can remove this per-resource-generated function -// and replace it with a single utility fetcher -// -// This function is auto-generated -func (s *Store) QueryComposeRecords( - ctx context.Context, - f composeType.RecordFilter, -) (_ []*composeType.Record, more bool, err error) { - var ( - ok bool - - set = make([]*composeType.Record, 0, DefaultSliceCapacity) - res *composeType.Record - aux *auxComposeRecord - rows *sql.Rows - count uint - expr, tExpr []goqu.Expression - - sortExpr []exp.OrderedExpression - ) - - if s.Filters.ComposeRecord != nil { - // extended filter set - tExpr, f, err = s.Filters.ComposeRecord(s, f) - } else { - // using generated filter - tExpr, f, err = ComposeRecordFilter(f) - } - - if err != nil { - err = fmt.Errorf("could generate filter expression for ComposeRecord: %w", err) - return - } - - expr = append(expr, tExpr...) - - // paging feature is enabled - if f.PageCursor != nil { - if tExpr, err = cursorWithSorting(f.PageCursor, s.sortableComposeRecordFields()); err != nil { - return - } else { - expr = append(expr, tExpr...) - } - } - - query := composeRecordSelectQuery(s.Dialect).Where(expr...) - - // sorting feature is enabled - if sortExpr, err = order(f.Sort, s.sortableComposeRecordFields()); err != nil { - err = fmt.Errorf("could generate order expression for ComposeRecord: %w", err) - return - } - - if len(sortExpr) > 0 { - query = query.Order(sortExpr...) - } - - if f.Limit > 0 { - query = query.Limit(f.Limit) - } - - rows, err = s.Query(ctx, query) - if err != nil { - err = fmt.Errorf("could not query ComposeRecord: %w", err) - return - } - - if err = rows.Err(); err != nil { - err = fmt.Errorf("could not query ComposeRecord: %w", err) - return - } - - defer func() { - closeError := rows.Close() - if err == nil { - // return error from close - err = closeError - } - }() - - for rows.Next() { - if err = rows.Err(); err != nil { - err = fmt.Errorf("could not query ComposeRecord: %w", err) - return - } - - aux = new(auxComposeRecord) - if err = aux.scan(rows); err != nil { - err = fmt.Errorf("could not scan rows for ComposeRecord: %w", err) - return - } - - count++ - if res, err = aux.decode(); err != nil { - err = fmt.Errorf("could not decode ComposeRecord: %w", err) - return - } - - // check fn set, call it and see if it passed the test - // if not, skip the item - if f.Check != nil { - if ok, err = f.Check(res); err != nil { - return - } else if !ok { - continue - } - } - - set = append(set, res) - } - - return set, f.Limit > 0 && count >= f.Limit, err - -} - -// LookupComposeRecordByID searches for compose record by ID -// -// It returns compose record even if deleted -// -// This function is auto-generated -func (s *Store) LookupComposeRecordByID(ctx context.Context, mod *composeType.Module, id uint64) (_ *composeType.Record, err error) { - var ( - rows *sql.Rows - aux = new(auxComposeRecord) - lookup = composeRecordSelectQuery(s.Dialect).Where( - goqu.I("id").Eq(id), - ).Limit(1) - ) - - rows, err = s.Query(ctx, lookup) - if err != nil { - return - } - - defer func() { - closeError := rows.Close() - if err == nil { - // return error from close - err = closeError - } - }() - - if err = rows.Err(); err != nil { - return - } - - if !rows.Next() { - return nil, store.ErrNotFound.Stack(1) - } - - if err = aux.scan(rows); err != nil { - return - } - - return aux.decode() -} - -// sortableComposeRecordFields returns all columns flagged as sortable -// -// With optional string arg, all columns are returned aliased -// -// This function is auto-generated -func (Store) sortableComposeRecordFields() map[string]string { - return map[string]string{ - "created_at": "created_at", - "createdat": "created_at", - "deleted_at": "deleted_at", - "deletedat": "deleted_at", - "id": "id", - "updated_at": "updated_at", - "updatedat": "updated_at", - } -} - -// collectComposeRecordCursorValues collects values from the given resource that and sets them to the cursor -// to be used for pagination -// -// Values that are collected must come from sortable, unique or primary columns/fields -// At least one of the collected columns must be flagged as unique, otherwise fn appends primary keys at the end -// -// Known issue: -// when collecting cursor values for query that sorts by unique column with partial index (ie: unique handle on -// undeleted items) -// -// This function is auto-generated -func (s *Store) collectComposeRecordCursorValues(res *composeType.Record, cc ...*filter.SortExpr) *filter.PagingCursor { - var ( - cur = &filter.PagingCursor{LThen: filter.SortExprSet(cc).Reversed()} - - hasUnique bool - - pkID bool - - collect = func(cc ...*filter.SortExpr) { - for _, c := range cc { - switch c.Column { - case "id": - cur.Set(c.Column, res.ID, c.Descending) - pkID = true - case "createdAt": - cur.Set(c.Column, res.CreatedAt, c.Descending) - case "updatedAt": - cur.Set(c.Column, res.UpdatedAt, c.Descending) - case "deletedAt": - cur.Set(c.Column, res.DeletedAt, c.Descending) - } - } - } - ) - - collect(cc...) - if !hasUnique || !pkID { - collect(&filter.SortExpr{Column: "id", Descending: false}) - } - - return cur - -} - -// checkComposeRecordConstraints performs lookups (on valid) resource to check if any of the values on unique fields -// already exists in the store -// -// Using built-in constraint checking would be more performant, but unfortunately we cannot rely -// on the full support (MySQL does not support conditional indexes) -// -// This function is auto-generated -func (s *Store) checkComposeRecordConstraints(ctx context.Context, res *composeType.Record) (err error) { - return nil -} - // CreateComposeRecordValue creates one or more rows in composeRecordValue collection // // This function is auto-generated diff --git a/store/interfaces.gen.go b/store/interfaces.gen.go index 767cdb83f..33b7e297b 100644 --- a/store/interfaces.gen.go +++ b/store/interfaces.gen.go @@ -17,7 +17,6 @@ import ( labelsType "github.com/cortezaproject/corteza-server/pkg/label/types" "github.com/cortezaproject/corteza-server/pkg/locale" rbacType "github.com/cortezaproject/corteza-server/pkg/rbac" - "github.com/cortezaproject/corteza-server/pkg/report" systemType "github.com/cortezaproject/corteza-server/system/types" "go.uber.org/zap" "golang.org/x/text/language" @@ -57,7 +56,6 @@ type ( ComposeModuleFields ComposeNamespaces ComposePages - ComposeRecords ComposeRecordValues Credentials DalConnections @@ -306,20 +304,6 @@ type ( ReorderComposePages(ctx context.Context, namespace_id uint64, parent_id uint64, page_ids []uint64) error } - ComposeRecords interface { - SearchComposeRecords(ctx context.Context, mod *composeType.Module, f composeType.RecordFilter) (composeType.RecordSet, composeType.RecordFilter, error) - CreateComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) error - UpdateComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) error - UpsertComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) error - DeleteComposeRecord(ctx context.Context, mod *composeType.Module, rr ...*composeType.Record) error - DeleteComposeRecordByID(ctx context.Context, mod *composeType.Module, id uint64) error - TruncateComposeRecords(ctx context.Context, mod *composeType.Module) error - LookupComposeRecordByID(ctx context.Context, mod *composeType.Module, id uint64) (*composeType.Record, error) - ComposeRecordReport(ctx context.Context, mod *composeType.Module, metrics string, dimensions string, filters string) ([]map[string]interface{}, error) - ComposeRecordDatasource(ctx context.Context, mod *composeType.Module, ld *report.LoadStepDefinition) (report.Datasource, error) - PartialComposeRecordValueUpdate(ctx context.Context, mod *composeType.Module, values ...*composeType.RecordValue) error - } - ComposeRecordValues interface { SearchComposeRecordValues(ctx context.Context, f composeType.RecordValueFilter) (composeType.RecordValueSet, composeType.RecordValueFilter, error) CreateComposeRecordValue(ctx context.Context, rr ...*composeType.RecordValue) error @@ -1806,85 +1790,6 @@ func ReorderComposePages(ctx context.Context, s ComposePages, namespace_id uint6 return s.ReorderComposePages(ctx, namespace_id, parent_id, page_ids) } -// SearchComposeRecords returns all matching ComposeRecords from store -// -// This function is auto-generated -func SearchComposeRecords(ctx context.Context, s ComposeRecords, mod *composeType.Module, f composeType.RecordFilter) (composeType.RecordSet, composeType.RecordFilter, error) { - return s.SearchComposeRecords(ctx, mod, f) -} - -// CreateComposeRecord creates one or more ComposeRecords in store -// -// This function is auto-generated -func CreateComposeRecord(ctx context.Context, s ComposeRecords, mod *composeType.Module, rr ...*composeType.Record) error { - return s.CreateComposeRecord(ctx, mod, rr...) -} - -// UpdateComposeRecord updates one or more (existing) ComposeRecords in store -// -// This function is auto-generated -func UpdateComposeRecord(ctx context.Context, s ComposeRecords, mod *composeType.Module, rr ...*composeType.Record) error { - return s.UpdateComposeRecord(ctx, mod, rr...) -} - -// UpsertComposeRecord creates new or updates existing one or more ComposeRecords in store -// -// This function is auto-generated -func UpsertComposeRecord(ctx context.Context, s ComposeRecords, mod *composeType.Module, rr ...*composeType.Record) error { - return s.UpsertComposeRecord(ctx, mod, rr...) -} - -// DeleteComposeRecord deletes one or more ComposeRecords from store -// -// This function is auto-generated -func DeleteComposeRecord(ctx context.Context, s ComposeRecords, mod *composeType.Module, rr ...*composeType.Record) error { - return s.DeleteComposeRecord(ctx, mod, rr...) -} - -// DeleteComposeRecordByID deletes one or more ComposeRecords from store -// -// This function is auto-generated -func DeleteComposeRecordByID(ctx context.Context, s ComposeRecords, mod *composeType.Module, id uint64) error { - return s.DeleteComposeRecordByID(ctx, mod, id) -} - -// TruncateComposeRecords Deletes all ComposeRecords from store -// -// This function is auto-generated -func TruncateComposeRecords(ctx context.Context, s ComposeRecords, mod *composeType.Module) error { - return s.TruncateComposeRecords(ctx, mod) -} - -// LookupComposeRecordByID searches for compose record by ID -// -// It returns compose record even if deleted -// -// This function is auto-generated -func LookupComposeRecordByID(ctx context.Context, s ComposeRecords, mod *composeType.Module, id uint64) (*composeType.Record, error) { - return s.LookupComposeRecordByID(ctx, mod, id) -} - -// ComposeRecordReport -// -// This function is auto-generated -func ComposeRecordReport(ctx context.Context, s ComposeRecords, mod *composeType.Module, metrics string, dimensions string, filters string) ([]map[string]interface{}, error) { - return s.ComposeRecordReport(ctx, mod, metrics, dimensions, filters) -} - -// ComposeRecordDatasource -// -// This function is auto-generated -func ComposeRecordDatasource(ctx context.Context, s ComposeRecords, mod *composeType.Module, ld *report.LoadStepDefinition) (report.Datasource, error) { - return s.ComposeRecordDatasource(ctx, mod, ld) -} - -// PartialComposeRecordValueUpdate -// -// This function is auto-generated -func PartialComposeRecordValueUpdate(ctx context.Context, s ComposeRecords, mod *composeType.Module, values ...*composeType.RecordValue) error { - return s.PartialComposeRecordValueUpdate(ctx, mod, values...) -} - // SearchComposeRecordValues returns all matching ComposeRecordValues from store // // This function is auto-generated diff --git a/store/tests/all_test.go b/store/tests/all_test.go index dc8fd493b..1d7960214 100644 --- a/store/tests/all_test.go +++ b/store/tests/all_test.go @@ -68,9 +68,6 @@ func testAllGenerated(t *testing.T, s store.Storer) { t.Run("composePage", func(t *testing.T) { testComposePages(t, s) }) - t.Run("composeRecord", func(t *testing.T) { - testComposeRecords(t, s) - }) t.Run("composeRecordValue", func(t *testing.T) { testComposeRecordValues(t, s) }) diff --git a/system/dalutils/connection.go b/system/dalutils/connection.go new file mode 100644 index 000000000..c7773391d --- /dev/null +++ b/system/dalutils/connection.go @@ -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 +} diff --git a/system/dalutils/sensitivity_level.go b/system/dalutils/sensitivity_level.go new file mode 100644 index 000000000..806d2e3f1 --- /dev/null +++ b/system/dalutils/sensitivity_level.go @@ -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 +} diff --git a/system/service/dal_connection.go b/system/service/dal_connection.go index 596b5caf1..2bfd28e80 100644 --- a/system/service/dal_connection.go +++ b/system/service/dal_connection.go @@ -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... +} diff --git a/system/service/dal_sensitivity_level.go b/system/service/dal_sensitivity_level.go index 5e591098c..98bc43b8e 100644 --- a/system/service/dal_sensitivity_level.go +++ b/system/service/dal_sensitivity_level.go @@ -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) { diff --git a/system/types/dal_connection.go b/system/types/dal_connection.go index aeffdfd9e..be0ff11e1 100644 --- a/system/types/dal_connection.go +++ b/system/types/dal_connection.go @@ -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