diff --git a/compose/crs/connection_wrap.go b/compose/crs/connection_wrap.go new file mode 100644 index 000000000..f5f3d993b --- /dev/null +++ b/compose/crs/connection_wrap.go @@ -0,0 +1,34 @@ +package crs + +import "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + +type ( + crsConnectionWrap struct { + id uint64 + dsn string + capabilities capabilities.Set + } +) + +// CRSConnectionWrap is a utility to define generic store connections +func CRSConnectionWrap(id uint64, dsn string, cc ...capabilities.Capability) crsConnectionWrap { + return crsConnectionWrap{ + id: id, + dsn: dsn, + capabilities: cc, + } +} + +// Receivers to conform to interface + +func (s crsConnectionWrap) ComposeRecordStoreID() uint64 { + return s.id +} + +func (s crsConnectionWrap) StoreDSN() string { + return s.dsn +} + +func (s crsConnectionWrap) Capabilities() capabilities.Set { + return s.capabilities +} diff --git a/compose/crs/crs.go b/compose/crs/crs.go new file mode 100644 index 000000000..57dcf89dc --- /dev/null +++ b/compose/crs/crs.go @@ -0,0 +1,173 @@ +package crs + +import ( + "context" + "fmt" + + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/cortezaproject/corteza-server/pkg/data" +) + +type ( + // the core struct that outlines the compose record store facility + composeRecordStore struct { + drivers []driver + stores map[uint64]*storeWrap + + // Indexed by corresponding storeID + models map[uint64]data.ModelSet + + primary *storeWrap + } + + storeWrap struct { + store connFunc + driver driver + dsn string + capabilities capabilities.Set + } + + connFunc func(ctx context.Context) (Store, error) + + crsDefiner interface { + ComposeRecordStoreID() uint64 + StoreDSN() string + Capabilities() capabilities.Set + } +) + +const ( + defaultStoreID uint64 = 0 +) + +// ComposeRecordStore initializes a fresh record store where the given store serves as the default +func ComposeRecordStore(ctx context.Context, primary crsDefiner, d driver, drivers ...driver) (*composeRecordStore, error) { + drivers = append([]driver{d}, drivers...) + + crs := &composeRecordStore{ + drivers: drivers, + stores: make(map[uint64]*storeWrap), + models: make(map[uint64]data.ModelSet), + primary: nil, + } + + d = crs.getDriver(primary) + if d == nil { + return nil, fmt.Errorf("could not add default store: no supported driver found") + } + + crs.primary = &storeWrap{ + store: storeConnWrap(d, primary), + driver: d, + dsn: primary.StoreDSN(), + capabilities: primary.Capabilities(), + } + + return crs, nil +} + +// AddStore registers the given store definitions as compose record stores +func (crs *composeRecordStore) AddStore(ctx context.Context, definers ...crsDefiner) error { + for _, definer := range definers { + if crs.stores[definer.ComposeRecordStoreID()] != nil { + return fmt.Errorf("can not add compose record store %d: already defined", definer.ComposeRecordStoreID()) + } + + d := crs.getDriver(definer) + if d == nil { + return fmt.Errorf("could not add store %d: no supported driver found", definer.ComposeRecordStoreID()) + } + + crs.stores[definer.ComposeRecordStoreID()] = &storeWrap{ + store: storeConnWrap(d, definer), + driver: d, + dsn: definer.StoreDSN(), + capabilities: definer.Capabilities(), + } + } + return nil +} + +// RemoveStore removes the given store definition as a compose record store +func (crs *composeRecordStore) RemoveStore(ctx context.Context, storeIDs ...uint64) (err error) { + for _, storeID := range storeIDs { + s := crs.stores[storeID] + if s == nil { + return fmt.Errorf("can not remove compose record store %d: store does not exist", storeID) + } + + // Any potential driver cleanup + err = s.driver.Close(ctx, s.dsn) + if err != nil { + return + } + + // Remove from registry + delete(crs.stores, storeID) + } + + return nil +} + +// --- + +func (crs *composeRecordStore) getModel(store uint64, ident string) *data.Model { + for _, model := range crs.models[store] { + if model.Ident == ident { + return model + } + } + + return nil +} + +// getStore returns a store for the given identifier/capabilities combination +func (crs *composeRecordStore) getStore(ctx context.Context, storeID uint64, cc ...capabilities.Capability) (store Store, can capabilities.Set, err error) { + err = func() error { + // get the requested store + var wrap *storeWrap + if storeID == defaultStoreID { + wrap = crs.primary + } else { + wrap = crs.stores[storeID] + } + if wrap == nil { + return fmt.Errorf("could not get store %d: store does not exist", storeID) + } + + // check if store supports requested capabilities + if !wrap.capabilities.IsSuperset(cc...) { + return fmt.Errorf("store does not support requested capabilities: %v", capabilities.Set(cc).Diff(wrap.capabilities)) + } + + store, err = wrap.store(ctx) + can = wrap.capabilities + return nil + }() + + if err != nil { + err = fmt.Errorf("could not connect to store %d: %v", storeID, err) + return + } + + return +} + +// getDriver returns a driver which can be used with the given store +func (crs *composeRecordStore) getDriver(def crsDefiner) driver { + for _, d := range crs.drivers { + if !d.Can(def.StoreDSN(), def.Capabilities()...) { + continue + } + + return d + } + + return nil +} + +func storeConnWrap(d driver, def crsDefiner) connFunc { + return func(ctx context.Context) (Store, error) { + return d.Store(ctx, def.StoreDSN()) + } +} diff --git a/compose/crs/data.go b/compose/crs/data.go new file mode 100644 index 000000000..e34d2d9d0 --- /dev/null +++ b/compose/crs/data.go @@ -0,0 +1,144 @@ +package crs + +import ( + "context" + "fmt" + + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/data" +) + +func (crs *composeRecordStore) ComposeRecordCreate(ctx context.Context, module *types.Module, records ...*types.Record) (err error) { + // Determine required capabilities + requiredCap := capabilities.CreateCapabilities(module.Store.Capabilities...) + + // Determine store + var s Store + if s, _, err = crs.getStore(ctx, module.Store.ComposeRecordStoreID, requiredCap...); err != nil { + return err + } + + // Prepare data + model := crs.lookupModel(module) + if model == nil { + return crs.modelNotFoundErr(module) + } + + aux := make([]Getter, len(records)) + for i, r := range records { + aux[i] = Getter(r) + } + return s.CreateRecords(ctx, model, aux...) +} + +func (crs *composeRecordStore) ComposeRecordSearch(ctx context.Context, module *types.Module, filter *types.RecordFilter) (records types.RecordSet, outFilter *types.RecordFilter, err error) { + // Determine requiredCap we'll need + requiredCap := capabilities.SearchCapabilities(module.Store.Capabilities...).Union(crs.recFilterCapabilities(filter)) + + // Connect to datasource + var s Store + var cc capabilities.Set + _ = cc + s, cc, err = crs.getStore(ctx, module.Store.ComposeRecordStoreID, requiredCap...) + if err != nil { + return + } + + // Prepare data + model := crs.lookupModel(module) + if model == nil { + return nil, nil, crs.modelNotFoundErr(module) + } + + loader, err := s.SearchRecords(ctx, model, nil) + if err != nil { + return + } + + limit := int(filter.Limit) + if limit == 0 { + limit = 10 + } + + auxCC := make([]Setter, limit) + for i := range auxCC { + auxCC[i] = &types.Record{} + } + + var ok bool + _ = ok + for loader.More() && len(records) < int(limit) { + _, err = loader.Load(model, auxCC) + if err != nil { + return + } + + auxRecords, err := crs.extractRecords(model, auxCC...) + if err != nil { + return nil, nil, err + } + + if !capabilities.AccessControlCapabilities().IsSubset(cc...) && filter.Check != nil { + for _, r := range auxRecords { + if r == nil { + continue + } + if ok, err = filter.Check(r); err != nil { + return nil, nil, err + } else if !ok { + continue + } + + records = append(records, r) + } + } else { + for _, r := range auxRecords { + if r == nil { + break + } + records = append(records, r) + } + } + } + + return +} + +// --- + +func (crs *composeRecordStore) recFilterCapabilities(f *types.RecordFilter) (out capabilities.Set) { + if f == nil { + return + } + 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 (crs composeRecordStore) modelNotFoundErr(module *types.Module) error { + return fmt.Errorf("cannot create records for module %d: module not registered to crs", module.ID) +} + +func (crs composeRecordStore) extractRecords(model *data.Model, gg ...Setter) (out types.RecordSet, err error) { + out = make(types.RecordSet, len(gg)) + for i, g := range gg { + out[i] = g.(*types.Record) + } + + return +} diff --git a/compose/crs/driver.go b/compose/crs/driver.go new file mode 100644 index 000000000..4035ea227 --- /dev/null +++ b/compose/crs/driver.go @@ -0,0 +1,68 @@ +package crs + +import ( + "context" + + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/cortezaproject/corteza-server/pkg/data" + "github.com/cortezaproject/corteza-server/pkg/expr" +) + +type ( + driver interface { + // Capabilities returns all of the capabilities the driver is able to support + Capabilities() capabilities.Set + // Can determines if the driver is able to handle the connection + Can(dsn string, capabilities ...capabilities.Capability) bool + + // Store connects and returns a store we can use + Store(ctx context.Context, dsn string) (Store, error) + // Close closes the connection to specified + Close(ctx context.Context, dsn string) error + } + + // Loader provides an interface for loading data from the underlying store + Loader interface { + More() bool + Load(*data.Model, []Setter) (coppied int, err error) + + // -1 means unknown + Total() int + Cursor() any + // ... do we need anything else here? + } + + // Store provides an interface which CRS uses to interact with the underlying database + Store interface { + // DML + CreateRecords(ctx context.Context, sch *data.Model, cc ...Getter) error + // @note providing a model here would allow us to control what attributes we return; would we find this useful? + SearchRecords(ctx context.Context, sch *data.Model, filter any) (Loader, error) + // Rest are same difference + // - LookupRecord + // - UpdateRecords + // - DeleteRecords + // - TruncateRecords + + // DDL + // Models returns all of the collections the given database already defines + Models(context.Context) (data.ModelSet, error) + // AddModel requests the driver to support the specified collections + AddModel(context.Context, *data.Model, ...*data.Model) error + // // RemoveModel requests the driver to remove support for the specified collections + // RemoveModel(context.Context, *data.Model, ...*data.Model) error + // AlterModel requests the driver to alter the general collectio parameters + AlterModel(ctx context.Context, old *data.Model, new *data.Model) error + // AlterModelAttribute requests the driver to alter the specified attribute of the given collection + AlterModelAttribute(ctx context.Context, sch *data.Model, old data.Attribute, new data.Attribute, trans ...func(*data.Model, data.Attribute, expr.TypedValue) (expr.TypedValue, bool, error)) error + } + + // probably somewhere else (on the db level? + Getter interface { + GetValue(name string, pos int) (expr.TypedValue, error) + } + + Setter interface { + SetValue(name string, pos int, value expr.TypedValue) error + } +) diff --git a/compose/crs/main_test.go b/compose/crs/main_test.go new file mode 100644 index 000000000..8fcb206ec --- /dev/null +++ b/compose/crs/main_test.go @@ -0,0 +1,94 @@ +package crs + +import ( + "time" + + "context" + "strings" + + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/data" + "github.com/cortezaproject/corteza-server/pkg/expr" + "github.com/cortezaproject/corteza-server/pkg/id" +) + +type ( + testDriver struct{} + + testStore map[string]*types.Record +) + +var ( + gStore = make(testStore) + + // wrapper around time.Now() that will aid service testing + now = func() *time.Time { + c := time.Now().Round(time.Second) + return &c + } + + // wrapper around nextID that will aid service testing + nextID = func() uint64 { + return id.Next() + } +) + +func (d testDriver) Capabilities() capabilities.Set { + return capabilities.FullCapabilities() +} + +func (d testDriver) Can(uri string, cc ...capabilities.Capability) bool { + return strings.HasPrefix(uri, "noop://") && capabilities.Set(cc).IsSubset(d.Capabilities()...) +} + +func (d testDriver) Store(ctx context.Context, uri string) (str Store, err error) { + return &testStore{}, nil +} + +func (d testDriver) Close(ctx context.Context, uri string) (err error) { + return nil +} + +// --- + +func (ds *testStore) CreateRecords(ctx context.Context, sch *data.Model, cc ...Getter) error { + return nil +} + +func (ds *testStore) SearchRecords(ctx context.Context, sch *data.Model, filter any) (Loader, error) { + return nil, nil +} + +func (ds *testStore) Models(context.Context) (data.ModelSet, error) { + return nil, nil +} + +func (ds *testStore) AddModel(context.Context, *data.Model, ...*data.Model) error { + return nil +} + +func (ds *testStore) AlterModel(ctx context.Context, old *data.Model, new *data.Model) error { + return nil +} + +func (ds *testStore) AlterModelAttribute(ctx context.Context, sch *data.Model, old data.Attribute, new data.Attribute, trans ...func(*data.Model, data.Attribute, expr.TypedValue) (expr.TypedValue, bool, error)) error { + return nil +} + +func initCRS(ctx context.Context) *composeRecordStore { + crs, _ := ComposeRecordStore(ctx, CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), testDriver{}) + crs.AddStore(ctx, CRSConnectionWrap(defaultExternalCRS, "noop://external", capabilities.FullCapabilities()...)) + return crs +} + +func defaultModule() *types.Module { + return &types.Module{ + ID: nextID(), + Store: types.CRSDef{ComposeRecordStoreID: defaultExternalCRS}, + Fields: types.ModuleFieldSet{{ + Name: "first_name", + Kind: "String", + }}, + } +} diff --git a/compose/crs/model.go b/compose/crs/model.go new file mode 100644 index 000000000..626e85871 --- /dev/null +++ b/compose/crs/model.go @@ -0,0 +1,473 @@ +package crs + +import ( + "context" + "fmt" + "strings" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/data" +) + +type ( + // cStore is a simplified interface so we can use the store.Storer to assert a valid schema + cStore interface { + SearchComposeModules(ctx context.Context, f types.ModuleFilter) (types.ModuleSet, types.ModuleFilter, error) + SearchComposeModuleFields(ctx context.Context, f types.ModuleFieldFilter) (types.ModuleFieldSet, types.ModuleFieldFilter, error) + SearchComposeNamespaces(ctx context.Context, f types.NamespaceFilter) (types.NamespaceSet, types.NamespaceFilter, error) + } +) + +const ( + // https://www.rfc-editor.org/errata/eid1690 + emailLength = 254 + + // Generally the upper most limit + urlLength = 2048 +) + +// ReloadModulesFromStore resets state based on the provided cStore +func (crs *composeRecordStore) ReloadModulesFromStore(ctx context.Context, cs cStore) (err error) { + modules, err := crs.loadModules(ctx, cs) + if err != nil { + return + } + + return crs.ReloadModules(ctx, modules...) +} + +// ReloadModulesFromStore resets state based on the provided set of modules +func (crs *composeRecordStore) ReloadModules(ctx context.Context, modules ...*types.Module) (err error) { + // Clear up the old ones + // @todo profile if manually removing nested pointers makes it faster + crs.models = make(map[uint64]data.ModelSet) + + return crs.AddModules(ctx, modules...) +} + +// AddModules adds new modules without affecting existing ones +func (crs *composeRecordStore) AddModules(ctx context.Context, modules ...*types.Module) (err error) { + models, err := crs.modulesToModel(modules...) + if err != nil { + return + } + + for storeID, models := range crs.modelByStore(models) { + if crs.stores[storeID] == nil { + return fmt.Errorf("can not add module to store %d: store does not exist", storeID) + } + + err = crs.addModel(ctx, storeID, models) + if err != nil { + return + } + } + + return +} + +// RemoveModules removes the specified modules +func (crs *composeRecordStore) RemoveModules(ctx context.Context, modules ...*types.Module) (err error) { + models, err := crs.modulesToModel(modules...) + if err != nil { + return + } + + // validation + for _, model := range models { + // Validate existence + old := crs.getModel(model.StoreID, model.Ident) + if old == nil { + return fmt.Errorf("cannot remove module %s: not registered", model.Ident) + } + + // Validate no leftover references + // @todo we can probably expand on this quitea bit + for _, registered := range crs.models { + refs := registered.FilterByReferenced(model) + if len(refs) > 0 { + return fmt.Errorf("cannot remove module %s: referenced by other modules", model.Ident) + } + } + } + + // Work + for _, model := range models { + oldModels := crs.models[model.StoreID] + crs.models[model.StoreID] = make(data.ModelSet, 0, len(oldModels)) + for _, o := range oldModels { + if o.Ident == model.Ident { + continue + } + + crs.models[model.StoreID] = append(crs.models[model.StoreID], o) + + } + + // @todo should the underlying store be notified about this? + } + + return nil +} + +// AlterModule updates the old module with the new one +func (crs *composeRecordStore) AlterModule(ctx context.Context, oldMod, newMod *types.Module) (err error) { + // validation + { + if oldMod.Store.ComposeRecordStoreID != newMod.Store.ComposeRecordStoreID { + return fmt.Errorf("cannot alter module stored in different record stores: old: %d, new: %d", oldMod.Store.ComposeRecordStoreID, newMod.Store.ComposeRecordStoreID) + } + } + + store, oldModel, err := crs.prepModuleDDL(ctx, oldMod) + if err != nil { + return + } + // store is same so we omit + _, newModel, err := crs.prepModuleDDL(ctx, newMod) + if err != nil { + return + } + + return store.AlterModel(ctx, oldModel, newModel) +} + +// @todo other ddl manipupations... + +// --- + +func (crs *composeRecordStore) prepModuleDDL(ctx context.Context, module *types.Module) (s Store, model *data.Model, err error) { + s, _, err = crs.getStore(ctx, module.Store.ComposeRecordStoreID) + if err != nil { + return + } + + models, err := crs.modulesToModel(module) + if err != nil { + return + } + model = models[0] + + return +} + +// modulesByStore maps the given module set by their CRS +func (crs *composeRecordStore) modulesByStore(modules types.ModuleSet) (out map[uint64]types.ModuleSet) { + out = make(map[uint64]types.ModuleSet) + for _, mod := range modules { + out[mod.Store.ComposeRecordStoreID] = append(out[mod.Store.ComposeRecordStoreID], mod) + } + + return +} + +// modelByStore maps the given models by their CRS +func (crs *composeRecordStore) modelByStore(models data.ModelSet) (out map[uint64]data.ModelSet) { + out = make(map[uint64]data.ModelSet) + + for _, model := range models { + out[model.StoreID] = append(out[model.StoreID], model) + } + + return +} + +// loadModules is a utility to load available modules with all their metadata included +func (crs *composeRecordStore) loadModules(ctx context.Context, cs cStore) (mm types.ModuleSet, err error) { + var ( + namespaces types.NamespaceSet + modules types.ModuleSet + fields types.ModuleFieldSet + ) + + namespaces, _, err = cs.SearchComposeNamespaces(ctx, types.NamespaceFilter{}) + if err != nil { + return + } + + for _, ns := range namespaces { + modules, _, err = cs.SearchComposeModules(ctx, types.ModuleFilter{ + NamespaceID: ns.ID, + }) + if err != nil { + return + } + + for _, mod := range modules { + fields, _, err = cs.SearchComposeModuleFields(ctx, types.ModuleFieldFilter{ + ModuleID: []uint64{mod.ID}, + }) + if err != nil { + return + } + + mod.Fields = append(mod.Fields, fields...) + } + } + + return +} + +func moduleFieldCodec(f *types.ModuleField) (strat data.StoreStrategy) { + // Defaulting to alias + strat = data.StoreCodecAlias{ + Ident: f.Name, + } + + switch { + case f.Encoding.EncodingStrategyAlias != nil: + strat = data.StoreCodecAlias{ + Ident: f.Encoding.EncodingStrategyAlias.Ident, + } + case f.Encoding.EncodingStrategyJSON != nil: + strat = data.StoreCodecJSON{ + Ident: f.Encoding.EncodingStrategyJSON.Ident, + Path: f.Encoding.EncodingStrategyJSON.Path, + } + } + + return +} + +// ---------- + +func (crs *composeRecordStore) lookupModel(module *types.Module) (out *data.Model) { + for _, model := range crs.models[module.Store.ComposeRecordStoreID] { + if model.ResourceID == module.ID { + return model + } + } + + return nil +} + +func (crs *composeRecordStore) modulesToModel(modules ...*types.Module) (out data.ModelSet, err error) { + refIndex := make(map[uint64]*data.Model) + out = make(data.ModelSet, 0, len(modules)) + + // Initial pass to get everything we can + for _, module := range modules { + model := crs.moduleModelInit(module) + refIndex[module.ID] = model + out = append(out, model) + } + + // Add stuff we already have + for _, models := range crs.models { + for _, model := range models { + refIndex[model.ResourceID] = model + } + } + + // Build up fields + for i, mod := range modules { + out[i].Attributes, err = crs.moduleModelAttributes(mod, refIndex) + if err != nil { + return + } + } + + return +} + +func (crs *composeRecordStore) moduleModelInit(mod *types.Module) (out *data.Model) { + return &data.Model{ + StoreID: mod.Store.ComposeRecordStoreID, + ResourceID: mod.ID, + ResourceType: types.ModuleResourceType, + + // @todo this isn't exactly this + Ident: mod.Handle, + Attributes: make(data.AttributeSet, len(mod.Fields)), + } +} + +func (crs *composeRecordStore) moduleModelAttributes(mod *types.Module, refIndex map[uint64]*data.Model) (out data.AttributeSet, err error) { + for _, f := range mod.Fields { + attr := &data.Attribute{ + Ident: f.Name, + MultiValue: f.Multi, + Store: moduleFieldCodec(f), + } + out = append(out, attr) + + switch strings.ToLower(f.Kind) { + case "bool": + attr.Type = data.TypeBoolean{} + case "datetime": + switch { + case f.IsDateOnly(): + attr.Type = data.TypeDate{} + case f.IsTimeOnly(): + attr.Type = data.TypeTime{} + default: + attr.Type = data.TypeTimestamp{} + } + case "email": + attr.Type = data.TypeText{Length: emailLength} + case "file": + attr.Type = data.TypeRef{ + RefModel: &data.Model{Ident: "attachments"}, + RefAttribute: &data.Attribute{Ident: "id"}, + } + case "number": + attr.Type = data.TypeNumber{ + Precision: f.Options.Precision(), + // Scale: , + } + case "record": + var refModel *data.Model + mRefID := f.Options.UInt64("moduleID") + if mRefID > 0 { + refModel = refIndex[mRefID] + } + + attr.Type = data.TypeRef{ + RefModel: refModel, + RefAttribute: &data.Attribute{ + Ident: "id", + }, + } + case "select": + attr.Type = data.TypeEnum{ + Values: f.SelectOptions(), + } + case "string": + attr.Type = data.TypeText{ + Length: 0, + } + case "url": + attr.Type = data.TypeText{ + Length: urlLength, + } + case "user": + attr.Type = data.TypeRef{ + // @todo... + + RefAttribute: &data.Attribute{ + Ident: "id", + }, + } + } + } + + // System attrs + out = append(out, + &data.Attribute{ + Ident: "recordID", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "recordID", + }), + Type: data.TypeID{}, + }, + &data.Attribute{ + Ident: "createdAt", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "createdAt", + }), + Type: data.TypeTimestamp{}, + }, + &data.Attribute{ + Ident: "updatedAt", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "updatedAt", + }), + Type: data.TypeTimestamp{}, + }, + &data.Attribute{ + Ident: "deletedAt", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "deletedAt", + }), + Type: data.TypeTimestamp{}, + }, + + &data.Attribute{ + Ident: "ownedBy", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "ownedBy", + }), + Type: data.TypeRef{ + RefAttribute: &data.Attribute{Ident: "id"}, + }, + }, + &data.Attribute{ + Ident: "createdBy", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "createdBy", + }), + Type: data.TypeRef{ + RefAttribute: &data.Attribute{Ident: "id"}, + }, + }, + &data.Attribute{ + Ident: "updatedBy", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "updatedBy", + }), + Type: data.TypeRef{ + RefAttribute: &data.Attribute{Ident: "id"}, + }, + }, + &data.Attribute{ + Ident: "deletedBy", + Store: moduleFieldCodec(&types.ModuleField{ + Name: "deletedBy", + }), + Type: data.TypeRef{ + RefAttribute: &data.Attribute{Ident: "id"}, + }, + }, + ) + + return +} + +func (crs *composeRecordStore) addModel(ctx context.Context, storeID uint64, models data.ModelSet) (err error) { + for _, model := range models { + existing := crs.getModel(storeID, model.Ident) + if existing != nil { + return fmt.Errorf("cannot add model %s to store %d: already exists", model.Ident, storeID) + } + + err = crs.addModelToStore(ctx, storeID, model) + if err != nil { + return + } + + crs.models[storeID] = append(crs.models[storeID], model) + } + + return +} + +func (crs *composeRecordStore) addModelToStore(ctx context.Context, storeID uint64, model *data.Model) (err error) { + s, _, err := crs.getStore(ctx, storeID) + if err != nil { + return err + } + + available, err := s.Models(ctx) + if err != nil { + return err + } + + // Check if already in there + if existing := available.FindByIdent(model.Ident); existing != nil { + // Assert validity + diff := existing.Diff(model) + if len(diff) > 0 { + return fmt.Errorf("model %s exists: model not compatible: %v", existing.Ident, diff) + } + + return nil + } + + // Try to add to store + err = s.AddModel(ctx, model) + if err != nil { + return + } + + return nil +} diff --git a/compose/crs/module_management_test.go b/compose/crs/module_management_test.go new file mode 100644 index 000000000..a9d5c07b4 --- /dev/null +++ b/compose/crs/module_management_test.go @@ -0,0 +1,77 @@ +package crs + +import ( + "context" + "testing" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/stretchr/testify/require" +) + +var ( + defaultExternalCRS uint64 = nextID() +) + +func TestLoadModules(t *testing.T) { + ctx := context.Background() + + t.Run("no supporting CRS", func(t *testing.T) { + crs := initCRS(ctx) + + err := crs.ReloadModules(ctx, &types.Module{ + ID: nextID(), + Store: types.CRSDef{ComposeRecordStoreID: 999}, + }) + require.Error(t, err) + }) + + t.Run("module added", func(t *testing.T) { + crs := initCRS(ctx) + + mod := defaultModule() + + err := crs.ReloadModules(ctx, mod) + require.NoError(t, err) + + model := crs.lookupModel(mod) + require.NotNil(t, model) + + // 8 sys, 1 custom + require.Len(t, model.Attributes, 8+1) + + require.Equal(t, defaultExternalCRS, model.StoreID) + }) + + t.Run("duplicate module added", func(t *testing.T) { + crs := initCRS(ctx) + + mod := defaultModule() + + err := crs.ReloadModules(ctx, mod, mod) + require.Error(t, err) + }) +} + +func TestRemoveModules(t *testing.T) { + ctx := context.Background() + + t.Run("module not found", func(t *testing.T) { + crs := initCRS(ctx) + + err := crs.RemoveModules(ctx, defaultModule()) + require.Error(t, err) + }) + + t.Run("module removed", func(t *testing.T) { + crs := initCRS(ctx) + mod := defaultModule() + + err := crs.ReloadModules(ctx, mod) + require.NoError(t, err) + + err = crs.RemoveModules(ctx, mod) + require.NoError(t, err) + + require.Nil(t, crs.lookupModel(mod)) + }) +} diff --git a/compose/crs/noop/noop.go b/compose/crs/noop/noop.go new file mode 100644 index 000000000..0bda091bb --- /dev/null +++ b/compose/crs/noop/noop.go @@ -0,0 +1,66 @@ +// This package is just so that I have an interface conforming driver + +package noop + +import ( + "context" + "strings" + + "github.com/cortezaproject/corteza-server/compose/crs" + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/data" + "github.com/cortezaproject/corteza-server/pkg/expr" +) + +type ( + noopDriver struct{} + + noopStore map[string]*types.Record +) + +var ( + gStore = make(noopStore) +) + +func Driver() *noopDriver { + return &noopDriver{} +} + +func (d noopDriver) Capabilities() capabilities.Set { + return capabilities.FullCapabilities() +} + +func (d noopDriver) Can(uri string, cc ...capabilities.Capability) bool { + return strings.HasPrefix(uri, "noop://") && capabilities.Set(cc).IsSubset(d.Capabilities()...) +} + +func (d noopDriver) Store(ctx context.Context, uri string) (str crs.Store, err error) { + return &noopStore{}, nil +} + +// --- + +func (ds *noopStore) CreateRecords(ctx context.Context, sch *data.Model, cc ...crs.Getter) error { + return nil +} + +func (ds *noopStore) SearchRecords(ctx context.Context, sch *data.Model, filter any) (crs.Loader, error) { + return nil, nil +} + +func (ds *noopStore) Models(context.Context) (data.ModelSet, error) { + return nil, nil +} + +func (ds *noopStore) AddModel(context.Context, *data.Model, ...*data.Model) error { + return nil +} + +func (ds *noopStore) AlterModel(ctx context.Context, old *data.Model, new *data.Model) error { + return nil +} + +func (ds *noopStore) AlterModelAttribute(ctx context.Context, sch *data.Model, old data.Attribute, new data.Attribute, trans ...func(*data.Model, data.Attribute, expr.TypedValue) (expr.TypedValue, bool, error)) error { + return nil +} diff --git a/compose/crs/store_management_test.go b/compose/crs/store_management_test.go new file mode 100644 index 000000000..c6177cd5d --- /dev/null +++ b/compose/crs/store_management_test.go @@ -0,0 +1,131 @@ +package crs + +import ( + "context" + "testing" + + "github.com/cortezaproject/corteza-server/compose/crs/capabilities" + "github.com/stretchr/testify/require" +) + +func TestInit(t *testing.T) { + ctx := context.Background() + + t.Run("invalid primary store", func(t *testing.T) { + _, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "invalid://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.Error(t, err) + }) + + t.Run("valid primary store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + require.Len(t, crs.stores, 0) + require.NotNil(t, crs.primary) + + }) +} + +func TestStoreAdd(t *testing.T) { + ctx := context.Background() + + t.Run("invalid external store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + sID := nextID() + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "invalid://external", capabilities.FullCapabilities()...)) + require.Error(t, err) + + require.NotNil(t, crs.primary) + require.Len(t, crs.stores, 0) + }) + + t.Run("valid external store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + sID := nextID() + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "noop://external", capabilities.FullCapabilities()...)) + require.NoError(t, err) + + require.NotNil(t, crs.primary) + require.Len(t, crs.stores, 1) + require.NotNil(t, crs.stores[sID]) + }) + + t.Run("duplicated external store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + sID := nextID() + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "noop://external", capabilities.FullCapabilities()...)) + require.NoError(t, err) + + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "noop://external", capabilities.FullCapabilities()...)) + require.Error(t, err) + + require.NotNil(t, crs.primary) + require.Len(t, crs.stores, 1) + require.NotNil(t, crs.stores[sID]) + }) +} + +func TestStoreRemove(t *testing.T) { + ctx := context.Background() + + t.Run("removing missing store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + sID := nextID() + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "noop://external", capabilities.FullCapabilities()...)) + require.NoError(t, err) + + err = crs.RemoveStore(ctx, 123) + require.Error(t, err) + require.Len(t, crs.stores, 1) + }) + + t.Run("removing store", func(t *testing.T) { + crs, err := ComposeRecordStore( + ctx, + CRSConnectionWrap(0, "noop://primary", capabilities.FullCapabilities()...), + testDriver{}, + ) + require.NoError(t, err) + + sID := nextID() + err = crs.AddStore(ctx, CRSConnectionWrap(sID, "noop://external", capabilities.FullCapabilities()...)) + require.NoError(t, err) + + err = crs.RemoveStore(ctx, sID) + require.NoError(t, err) + require.Len(t, crs.stores, 0) + require.NotNil(t, crs.primary) + }) +} diff --git a/compose/types/module.go b/compose/types/module.go index 76a41ada5..cb209becc 100644 --- a/compose/types/module.go +++ b/compose/types/module.go @@ -12,12 +12,23 @@ import ( ) type ( + CRSDef struct { + ComposeRecordStoreID uint64 + Capabilities capabilities.Set + + Partitioned bool + PartitionFormat string + PartitionBase []any + } + Module struct { ID uint64 `json:"moduleID,string"` Handle string `json:"handle"` Meta types.JSONText `json:"meta"` Fields ModuleFieldSet `json:"fields"` + Store CRSDef + Labels map[string]string `json:"labels,omitempty"` NamespaceID uint64 `json:"namespaceID,string"` diff --git a/compose/types/module_field.go b/compose/types/module_field.go index bf15ab3aa..d6028e7d9 100644 --- a/compose/types/module_field.go +++ b/compose/types/module_field.go @@ -57,6 +57,50 @@ var ( _ sort.Interface = &ModuleFieldSet{} ) +func (f *ModuleField) SelectOptions() (out []string) { + if f.Kind != "Select" { + return + } + + var ( + options, has = f.Options["options"] + ) + + if !has { + return + } + + switch oo := options.(type) { + case []string: + out = oo + case []interface{}: + for _, o := range oo { + switch c := o.(type) { + case string: + out = append(out, c) + case map[string]string: + if value, has := c["value"]; has { + out = append(out, value) + } + case map[string]interface{}: + if value, has := c["value"]; has { + if value, ok := value.(string); ok { + out = append(out, value) + } + } + case ModuleFieldOptions: + if value, has := c["value"]; has { + if value, ok := value.(string); ok { + out = append(out, value) + } + } + } + } + } + + return +} + func (f *ModuleField) decodeTranslationsExpressionValidatorValidatorIDError(tt locale.ResourceTranslationIndex) { var aux *locale.ResourceTranslation