From 31c672e2e42c28f4c90e23533c6e016f551dc754 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Toma=C5=BE=20Jerman?= Date: Tue, 17 Nov 2020 16:01:08 +0100 Subject: [PATCH] Implement base store encoder states --- pkg/envoy/store/application.go | 179 +++++++++++++++ pkg/envoy/store/compose_chart.go | 225 +++++++++++++++++++ pkg/envoy/store/compose_module.go | 324 +++++++++++++++++++++++++++ pkg/envoy/store/compose_namespace.go | 172 ++++++++++++++ pkg/envoy/store/compose_page.go | 229 +++++++++++++++++++ pkg/envoy/store/compose_record.go | 206 +++++++++++++++++ pkg/envoy/store/encoder.go | 149 ++++++++++++ pkg/envoy/store/rbac_rule.go | 131 +++++++++++ pkg/envoy/store/role.go | 166 ++++++++++++++ pkg/envoy/store/settings.go | 90 ++++++++ pkg/envoy/store/user.go | 181 +++++++++++++++ pkg/envoy/store/util.go | 54 +++++ 12 files changed, 2106 insertions(+) create mode 100644 pkg/envoy/store/application.go create mode 100644 pkg/envoy/store/compose_chart.go create mode 100644 pkg/envoy/store/compose_module.go create mode 100644 pkg/envoy/store/compose_namespace.go create mode 100644 pkg/envoy/store/compose_page.go create mode 100644 pkg/envoy/store/compose_record.go create mode 100644 pkg/envoy/store/encoder.go create mode 100644 pkg/envoy/store/rbac_rule.go create mode 100644 pkg/envoy/store/role.go create mode 100644 pkg/envoy/store/settings.go create mode 100644 pkg/envoy/store/user.go create mode 100644 pkg/envoy/store/util.go diff --git a/pkg/envoy/store/application.go b/pkg/envoy/store/application.go new file mode 100644 index 000000000..908bf1b94 --- /dev/null +++ b/pkg/envoy/store/application.go @@ -0,0 +1,179 @@ +package store + +import ( + "context" + "time" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" + "github.com/cortezaproject/corteza-server/system/types" +) + +type ( + applicationState struct { + cfg *EncoderConfig + + res *resource.Application + app *types.Application + } +) + +func NewApplicationState(res *resource.Application, cfg *EncoderConfig) resourceState { + return &applicationState{ + cfg: cfg, + res: res, + } +} + +func (n *applicationState) Prepare(ctx context.Context, s store.Storer, rs *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + if n.res.Res.Unify == nil { + n.res.Res.Unify = &types.ApplicationUnify{} + } + + // Get the existing app + n.app, err = findApplicationS(ctx, s, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.app != nil { + n.res.Res.ID = n.app.ID + } + return nil +} + +// Encode encodes the given application +func (n *applicationState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.app != nil && n.app.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.app.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Create fresh application + if !exists { + return store.CreateApplication(ctx, s, res) + } + + // Update existing application + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeApplications(n.app, res) + + case MergeRight: + res = mergeApplications(res, n.app) + } + + err = store.UpdateApplication(ctx, s, res) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeApplications merges b into a, prioritising a +func mergeApplications(a, b *types.Application) *types.Application { + c := *a + + if c.Name == "" { + c.Name = b.Name + } + + // I'll just compare the entire struct for now + if c.Unify == nil || *c.Unify == (types.ApplicationUnify{}) { + c.Unify = b.Unify + } + + return &c +} + +// findApplicationRS looks for the app in the resources & the store +// +// Provided resources are prioritized. +func findApplicationRS(ctx context.Context, s store.Storer, rr resource.InterfaceSet, ii resource.Identifiers) (ap *types.Application, err error) { + ap = findApplicationR(rr, ii) + if ap != nil { + return ap, nil + } + + return findApplicationS(ctx, s, makeGenericFilter(ii)) +} + +// findApplicationS looks for the app in the store +func findApplicationS(ctx context.Context, s store.Storer, gf genericFilter) (ap *types.Application, err error) { + if gf.id > 0 { + ap, err = store.LookupApplicationByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if ap != nil { + return + } + } + + q := gf.handle + if q == "" { + q = gf.name + } + + if q != "" { + var aa types.ApplicationSet + aa, _, err = store.SearchApplications(ctx, s, types.ApplicationFilter{Name: q}) + if err != nil && err != store.ErrNotFound { + return nil, err + } + if len(aa) > 0 { + ap = aa[0] + return + } + } + + return nil, nil +} + +// findApplicationR looks for the app in the resource set +func findApplicationR(rr resource.InterfaceSet, ii resource.Identifiers) (ap *types.Application) { + var apRes *resource.Application + var ok bool + + rr.Walk(func(r resource.Interface) error { + apRes, ok = r.(*resource.Application) + if !ok { + return nil + } + + if !apRes.Identifiers().HasAny(r.Identifiers()) { + apRes = nil + } + + return nil + }) + + // Found it + if apRes != nil { + return apRes.Res + } + + return nil +} diff --git a/pkg/envoy/store/compose_chart.go b/pkg/envoy/store/compose_chart.go new file mode 100644 index 000000000..6376f2d9d --- /dev/null +++ b/pkg/envoy/store/compose_chart.go @@ -0,0 +1,225 @@ +package store + +import ( + "context" + "errors" + "time" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + composeChartState struct { + cfg *EncoderConfig + + res *resource.ComposeChart + chr *types.Chart + + relNS *types.Namespace + relMods types.ModuleSet + } +) + +func NewComposeChartState(res *resource.ComposeChart, cfg *EncoderConfig) resourceState { + return &composeChartState{ + cfg: cfg, + + res: res, + } +} + +func (n *composeChartState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + // Get relate namespace + n.relNS, err = findComposeNamespaceRS(ctx, s, state.ParentResources, n.res.NsRef.Identifiers) + if err != nil { + return err + } + if n.relNS == nil { + return errors.New("@todo couldn't resolve namespace") + } + + // Can't do anything else, since the NS doesn't yet exist + if n.relNS.ID <= 0 { + return nil + } + + // Get related modules + n.relMods = make(types.ModuleSet, len(n.res.ModRef)) + for i, mRef := range n.res.ModRef { + n.relMods[i], err = findComposeModuleRS(ctx, s, n.relNS.ID, state.ParentResources, mRef.Identifiers) + if err != nil { + return err + } + } + + // Try to get the original chart + n.chr, err = findComposeChartS(ctx, s, n.relNS.ID, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.chr != nil { + n.res.Res.ID = n.chr.ID + } + return nil +} + +func (n *composeChartState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.chr != nil && n.chr.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.chr.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Namespace + res.NamespaceID = n.relNS.ID + if res.NamespaceID <= 0 { + ns := findComposeNamespaceR(state.ParentResources, n.res.NsRef.Identifiers) + res.NamespaceID = ns.ID + } + if res.NamespaceID <= 0 { + return errors.New("[chart] couldn't find related namespace; @todo error") + } + + // Report modules + for i, r := range res.Config.Reports { + mod := n.relMods[i] + if mod == nil { + mod = findComposeModuleR(state.ParentResources, n.res.ModRef[i].Identifiers) + } + if mod == nil || mod.ID <= 0 { + return errors.New("[chart] couldn't find related report module; @todo error") + } + + r.ModuleID = mod.ID + } + + // Create fresh chart + if !exists { + return store.CreateComposeChart(ctx, s, res) + } + + // Update existing chart + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeComposeChart(n.chr, res) + + case MergeRight: + res = mergeComposeChart(res, n.chr) + } + + err = store.UpdateComposeChart(ctx, s, res) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeComposeChart merges b into a, prioritising a +func mergeComposeChart(a, b *types.Chart) *types.Chart { + c := *a + + if c.Handle == "" { + c.Handle = b.Handle + } + if c.Name == "" { + c.Name = b.Name + } + + if len(c.Config.Reports) <= 0 { + c.Config.Reports = b.Config.Reports + } + if c.Config.ColorScheme == "" { + c.Config.ColorScheme = b.Config.ColorScheme + } + + return &c +} + +// findComposeChartRS looks for the chart in the resources & the store +// +// Provided resources are prioritized. +func findComposeChartRS(ctx context.Context, s store.Storer, nsID uint64, rr resource.InterfaceSet, ii resource.Identifiers) (ch *types.Chart, err error) { + ch = findComposeChartR(rr, ii) + if ch != nil { + return ch, nil + } + + // Go in the store + return findComposeChartS(ctx, s, nsID, makeGenericFilter(ii)) +} + +// findComposeChartS looks for the chart in the store +func findComposeChartS(ctx context.Context, s store.Storer, nsID uint64, gf genericFilter) (ch *types.Chart, err error) { + if gf.id > 0 { + ch, err = store.LookupComposeChartByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if ch != nil { + return + } + } + + if gf.handle != "" { + ch, err = store.LookupComposeChartByNamespaceIDHandle(ctx, s, nsID, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if ch != nil { + return + } + } + + return nil, nil +} + +// findComposeChartR looks for the chart in the resources +func findComposeChartR(rr resource.InterfaceSet, ii resource.Identifiers) (ch *types.Chart) { + var chRes *resource.ComposeChart + var ok bool + + rr.Walk(func(r resource.Interface) error { + chRes, ok = r.(*resource.ComposeChart) + if !ok { + return nil + } + + if !chRes.Identifiers().HasAny(r.Identifiers()) { + chRes = nil + } + return nil + }) + + // Found it + if chRes != nil { + return chRes.Res + } + + return nil +} diff --git a/pkg/envoy/store/compose_module.go b/pkg/envoy/store/compose_module.go new file mode 100644 index 000000000..e461aaa75 --- /dev/null +++ b/pkg/envoy/store/compose_module.go @@ -0,0 +1,324 @@ +package store + +import ( + "context" + "errors" + "time" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + composeModuleState struct { + cfg *EncoderConfig + + res *resource.ComposeModule + mod *types.Module + + relNS *types.Namespace + recFields map[string]uint64 + } +) + +func NewComposeModuleState(res *resource.ComposeModule, cfg *EncoderConfig) resourceState { + return &composeModuleState{ + cfg: cfg, + + res: res, + } +} + +func (n *composeModuleState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + // Get relate namespace + n.relNS, err = findComposeNamespaceRS(ctx, s, state.ParentResources, n.res.NsRef.Identifiers) + if err != nil { + return err + } + if n.relNS == nil { + return errors.New("@todo couldn't resolve namespace") + } + + // Can't do anything else, since the NS doesn't yet exist + if n.relNS.ID <= 0 { + return nil + } + + // Get related record field modules + for _, r := range n.res.ModRef { + mod, err := findComposeModuleRS(ctx, s, n.relNS.ID, state.ParentResources, r.Identifiers) + if err != nil { + return err + } else if mod == nil { + return errors.New("[mod prepare] couldn't find related module; @todo error") + } + + for i := range r.Identifiers { + n.recFields[i] = mod.ID + } + } + + // Try to get the original module + n.mod, err = findComposeModuleS(ctx, s, n.relNS.ID, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + // Nothing else to do + if n.mod == nil { + return nil + } + + // Get the original module fields + // These are used later for some merging logic + n.mod.Fields, err = findComposeModuleFieldsS(ctx, s, n.mod) + if err != nil { + return err + } + + if n.mod != nil { + n.res.Res.ID = n.mod.ID + } + return nil +} + +func (n *composeModuleState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.mod != nil && n.mod.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.mod.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + if state.Conflicting { + return nil + } + + // Namespace + res.NamespaceID = n.relNS.ID + if res.NamespaceID <= 0 { + ns := findComposeNamespaceR(state.ParentResources, n.res.NsRef.Identifiers) + res.NamespaceID = ns.ID + } + + if res.NamespaceID <= 0 { + return errors.New("[module] couldn't find related namespace; @todo error") + } + + // Fields + for i, f := range res.Fields { + f.ID = res.ID + f.ModuleID = res.ID + f.Place = i + f.DeletedAt = nil + f.CreatedAt = time.Now() + + if f.Kind == "Record" { + refM := f.Options.String("module") + mID := n.recFields[refM] + if mID <= 0 { + mod := findComposeModuleR(state.ParentResources, resource.MakeIdentifiers(refM)) + if mod == nil || mod.ID <= 0 { + return errors.New("[module field] couldn't find related module; @todo error") + } + mID = mod.ID + } + + f.Options["module"] = mID + } + } + + // Create a fresh module + if !exists { + err = store.CreateComposeModule(ctx, s, res) + if err != nil { + return err + } + + err = store.CreateComposeModuleField(ctx, s, res.Fields...) + if err != nil { + return err + } + + return nil + } + + // Update existing module + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeComposeModule(n.mod, res) + res.Fields = mergeComposeModuleFields(n.mod.Fields, res.Fields) + + case MergeRight: + res = mergeComposeModule(res, n.mod) + res.Fields = mergeComposeModuleFields(res.Fields, n.mod.Fields) + } + + err = store.UpdateComposeModule(ctx, s, res) + if err != nil { + return err + } + + err = store.DeleteComposeModuleField(ctx, s, n.mod.Fields...) + if err != nil { + return err + } + err = store.CreateComposeModuleField(ctx, s, res.Fields...) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeComposeModuleFields merges b into a, prioritising a +func mergeComposeModuleFields(a, b types.ModuleFieldSet) types.ModuleFieldSet { + ff := a.Clone() + missing := make(types.ModuleFieldSet, 0, len(b)) + + for _, fb := range b { + for _, fa := range ff { + if fb.ID != fa.ID && fb.Name != fa.Name { + continue + } + + if fa.Kind == "" { + fa.Kind = fb.Kind + } + if fa.Name == "" { + fa.Name = fb.Name + } + if fa.Label == "" { + fa.Label = fb.Label + } + if fa.Options == nil { + fa.Options = fb.Options + } + if fa.DefaultValue == nil || len(fa.DefaultValue) <= 0 { + fa.DefaultValue = fb.DefaultValue + } + + goto out + } + missing = append(missing, fb) + out: + } + + ff = append(ff, missing...) + + return ff +} + +// mergeComposeModule merges b into a, prioritising a +func mergeComposeModule(a, b *types.Module) *types.Module { + c := a.Clone() + + if c.Handle == "" { + c.Handle = b.Handle + } + if c.Name == "" { + c.Name = b.Name + } + + // I'll just compare the entire struct for now + if c.Meta == nil { + c.Meta = b.Meta + } + + return c +} + +// findComposeModuleRS looks for the chart in the resources & the store +// +// Provided resources are prioritized. +func findComposeModuleRS(ctx context.Context, s store.Storer, nsID uint64, rr resource.InterfaceSet, ii resource.Identifiers) (mod *types.Module, err error) { + mod = findComposeModuleR(rr, ii) + if mod != nil { + return mod, nil + } + + // Go in the store + return findComposeModuleS(ctx, s, nsID, makeGenericFilter(ii)) +} + +// findComposeModuleS looks for the module in the store +func findComposeModuleS(ctx context.Context, s store.Storer, nsID uint64, gf genericFilter) (mod *types.Module, err error) { + if gf.id > 0 { + mod, err = store.LookupComposeModuleByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if mod != nil { + return + } + } + + if gf.handle != "" { + mod, err = store.LookupComposeModuleByNamespaceIDHandle(ctx, s, nsID, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if mod != nil { + return + } + } + + return nil, nil +} + +// findComposeModuleR looks for the module in the store +func findComposeModuleR(rr resource.InterfaceSet, ii resource.Identifiers) (ns *types.Module) { + var modRes *resource.ComposeModule + var ok bool + + rr.Walk(func(r resource.Interface) error { + modRes, ok = r.(*resource.ComposeModule) + if !ok { + return nil + } + + if !modRes.Identifiers().HasAny(r.Identifiers()) { + modRes = nil + } + return nil + }) + + // Found it + if modRes != nil { + return modRes.Res + } + return nil +} + +// findComposeModuleFieldsS looks for the module fields in the store +func findComposeModuleFieldsS(ctx context.Context, s store.Storer, mod *types.Module) (types.ModuleFieldSet, error) { + if mod.ID <= 0 { + return mod.Fields, nil + } + + ff, _, err := store.SearchComposeModuleFields(ctx, s, types.ModuleFieldFilter{ + ModuleID: []uint64{mod.ID}, + }) + if err != nil { + return nil, err + } + + return ff, nil +} diff --git a/pkg/envoy/store/compose_namespace.go b/pkg/envoy/store/compose_namespace.go new file mode 100644 index 000000000..055d209b5 --- /dev/null +++ b/pkg/envoy/store/compose_namespace.go @@ -0,0 +1,172 @@ +package store + +import ( + "context" + "time" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + composeNamespaceState struct { + cfg *EncoderConfig + + res *resource.ComposeNamespace + ns *types.Namespace + } +) + +func NewComposeNamespaceState(res *resource.ComposeNamespace, cfg *EncoderConfig) resourceState { + return &composeNamespaceState{ + cfg: cfg, + + res: res, + } +} + +func (n *composeNamespaceState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + // Try to get the original chart + n.ns, err = findComposeNamespaceS(ctx, s, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.ns != nil { + n.res.Res.ID = n.ns.ID + } + return nil +} + +func (n *composeNamespaceState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.ns != nil && n.ns.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.ns.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Create a fresh namespace + if !exists { + return store.CreateComposeNamespace(ctx, s, res) + } + + // Update existing namespace + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeComposeNamespaces(n.ns, res) + + case MergeRight: + res = mergeComposeNamespaces(res, n.ns) + } + + err = store.UpdateComposeNamespace(ctx, s, res) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeComposeNamespaces merges b into a, prioritising a +func mergeComposeNamespaces(a, b *types.Namespace) *types.Namespace { + c := a.Clone() + + if c.Name == "" { + c.Name = b.Name + } + if c.Slug == "" { + c.Slug = b.Slug + } + + // I'll just compare the entire struct for now + if c.Meta == (types.NamespaceMeta{}) { + c.Meta = b.Meta + } + + return c +} + +// findComposeNamespaceRS looks for the namespace in the resources & the store +// +// Provided resources are prioritized. +func findComposeNamespaceRS(ctx context.Context, s store.Storer, rr resource.InterfaceSet, ii resource.Identifiers) (ns *types.Namespace, err error) { + ns = findComposeNamespaceR(rr, ii) + if ns != nil { + return ns, nil + } + + return findComposeNamespaceS(ctx, s, makeGenericFilter(ii)) +} + +// findComposeNamespaceS looks for the namespace in the store +func findComposeNamespaceS(ctx context.Context, s store.Storer, gf genericFilter) (ns *types.Namespace, err error) { + if gf.id > 0 { + ns, err = store.LookupComposeNamespaceByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if ns != nil { + return + } + } + + if gf.handle != "" { + ns, err = store.LookupComposeNamespaceBySlug(ctx, s, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if ns != nil { + return + } + } + + return nil, nil +} + +// findComposeNamespaceR looks for the namespace in the resources +func findComposeNamespaceR(rr resource.InterfaceSet, ii resource.Identifiers) (ns *types.Namespace) { + var nsRes *resource.ComposeNamespace + var ok bool + + rr.Walk(func(r resource.Interface) error { + nsRes, ok = r.(*resource.ComposeNamespace) + if !ok { + return nil + } + + if !nsRes.Identifiers().HasAny(r.Identifiers()) { + nsRes = nil + } + return nil + }) + + // Found it + if nsRes != nil { + return nsRes.Res + } + + return nil +} diff --git a/pkg/envoy/store/compose_page.go b/pkg/envoy/store/compose_page.go new file mode 100644 index 000000000..92c95dfcf --- /dev/null +++ b/pkg/envoy/store/compose_page.go @@ -0,0 +1,229 @@ +package store + +import ( + "context" + "errors" + "time" + + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + composePageState struct { + cfg *EncoderConfig + + res *resource.ComposePage + pg *types.Page + + relNS *types.Namespace + relMod *types.Module + } +) + +func NewComposePageState(res *resource.ComposePage, cfg *EncoderConfig) resourceState { + return &composePageState{ + cfg: cfg, + + res: res, + } +} + +func (n *composePageState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + // Get relate namespace + n.relNS, err = findComposeNamespaceRS(ctx, s, state.ParentResources, n.res.NsRef.Identifiers) + if err != nil { + return err + } + if n.relNS == nil { + return errors.New("@todo couldn't resolve namespace") + } + + // Can't do anything else, since the NS doesn't yet exist + if n.relNS.ID <= 0 { + return nil + } + + // Get related module + // If this isn't a record page, there is no related module + if n.res.ModRef != nil { + n.relMod, err = findComposeModuleRS(ctx, s, n.relNS.ID, state.ParentResources, n.res.ModRef.Identifiers) + if err != nil { + return err + } + if n.relMod == nil { + return errors.New("@todo couldn't resolve module") + } + } + + // Try to get the original page + n.pg, err = findComposePageS(ctx, s, n.relNS.ID, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.pg != nil { + n.res.Res.ID = n.pg.ID + } + return nil +} + +func (n *composePageState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.pg != nil && n.pg.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.pg.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Namespace + res.NamespaceID = n.relNS.ID + if res.NamespaceID <= 0 { + ns := findComposeNamespaceR(state.ParentResources, n.res.NsRef.Identifiers) + res.NamespaceID = ns.ID + } + + if res.NamespaceID <= 0 { + return errors.New("[chart] couldn't find related namespace; @todo error") + } + + // Module? + if n.res.ModRef != nil { + res.ModuleID = n.relMod.ID + if res.ModuleID <= 0 { + mod := findComposeModuleR(state.ParentResources, n.res.ModRef.Identifiers) + res.ModuleID = mod.ID + } + } + + // Create a fresh page + if !exists { + return store.CreateComposePage(ctx, s, res) + } + + // Update existing page + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeComposePage(n.pg, res) + + case MergeRight: + res = mergeComposePage(res, n.pg) + } + + err = store.UpdateComposePage(ctx, s, res) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeComposePage merges b into a, prioritising a +func mergeComposePage(a, b *types.Page) *types.Page { + c := a.Clone() + + if c.SelfID <= 0 { + c.SelfID = b.SelfID + } + if c.Handle == "" { + c.Handle = b.Handle + } + if c.Title == "" { + c.Title = b.Title + } + if c.Description == "" { + c.Description = b.Description + } + if len(c.Blocks) <= 0 { + c.Blocks = b.Blocks + } + if len(c.Children) <= 0 { + c.Children = b.Children + } + + return c +} + +// findComposePageRS looks for the page in the resources & the store +// +// Provided resources are prioritized. +func findComposePageRS(ctx context.Context, s store.Storer, nsID uint64, rr resource.InterfaceSet, ii resource.Identifiers) (pg *types.Page, err error) { + pg = findComposePageR(ctx, rr, ii) + if pg != nil { + return pg, nil + } + + // Go in the store + return findComposePageS(ctx, s, nsID, makeGenericFilter(ii)) +} + +// findComposePageS looks for the page in the store +func findComposePageS(ctx context.Context, s store.Storer, nsID uint64, gf genericFilter) (pg *types.Page, err error) { + if gf.id > 0 { + pg, err = store.LookupComposePageByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if pg != nil { + return + } + } + + if gf.handle != "" { + pg, err = store.LookupComposePageByNamespaceIDHandle(ctx, s, nsID, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if pg != nil { + return + } + } + + return nil, nil +} + +// findComposePageR looks for the page in the resources +func findComposePageR(ctx context.Context, rr resource.InterfaceSet, ii resource.Identifiers) (pg *types.Page) { + var pgRes *resource.ComposePage + var ok bool + + rr.Walk(func(r resource.Interface) error { + pgRes, ok = r.(*resource.ComposePage) + if !ok { + return nil + } + + if !pgRes.Identifiers().HasAny(r.Identifiers()) { + pgRes = nil + } + return nil + }) + + // Found it + if pgRes != nil { + return pgRes.Res + } + return nil +} diff --git a/pkg/envoy/store/compose_record.go b/pkg/envoy/store/compose_record.go new file mode 100644 index 000000000..def24921d --- /dev/null +++ b/pkg/envoy/store/compose_record.go @@ -0,0 +1,206 @@ +package store + +import ( + "context" + "errors" + "strconv" + "time" + + "github.com/cortezaproject/corteza-server/compose/service/values" + "github.com/cortezaproject/corteza-server/compose/types" + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + composeRecordState struct { + cfg *EncoderConfig + + res *resource.ComposeRecord + + relNS *types.Namespace + relMod *types.Module + } +) + +var ( + rvSanitizer = values.Sanitizer() + rvValidator = values.Validator() +) + +func NewComposeRecordState(res *resource.ComposeRecord, cfg *EncoderConfig) resourceState { + return &composeRecordState{ + cfg: cfg, + + res: res, + } +} + +func (n *composeRecordState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Get related namespace + n.relNS, err = findComposeNamespaceRS(ctx, s, state.ParentResources, n.res.NsRef.Identifiers) + if err != nil { + return err + } + if n.relNS == nil { + return errors.New("@todo couldn't resolve namespace") + } + + n.relMod = findComposeModuleR(state.ParentResources, n.res.ModRef.Identifiers) + if n.relMod == nil && n.relNS.ID > 0 { + n.relMod, err = findComposeModuleS(ctx, s, n.relNS.ID, makeGenericFilter(n.res.ModRef.Identifiers)) + if err != nil { + return err + } + if n.relMod == nil { + return errors.New("@todo couldn't resolve module") + } + + // Preload existing fields + n.relMod.Fields, err = findComposeModuleFieldsS(ctx, s, n.relMod) + if err != nil { + return err + } + } + + // Add missing refs + for _, f := range n.relMod.Fields { + switch f.Kind { + case "Record": + refM := f.Options.String("module") + if refM != "" && refM != "0" { + // Make a reference with that module's records + n.res.AddRef(resource.COMPOSE_RECORD_RESOURCE_TYPE, refM) + } + } + } + + // Can't do anything else, since the NS doesn't yet exist + if n.relNS.ID <= 0 { + return nil + } + + // Try to get existing records + // + // @todo handle large amounts of + for rID := range n.res.IDMap { + var r *types.Record + // @todo support for labels + if refy.MatchString(rID) { + id, _ := strconv.ParseUint(rID, 10, 64) + r, err = store.LookupComposeRecordByID(ctx, s, n.relMod, id) + if err == store.ErrNotFound { + continue + } else if err != nil { + return err + } + if r != nil { + n.res.RecMap[rID] = r + } + } else { + continue + } + + n.res.RecMap[rID] = r + } + + return nil +} + +func (n *composeRecordState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Namespace + nsID := n.relNS.ID + if nsID <= 0 { + ns := findComposeNamespaceR(state.ParentResources, n.res.NsRef.Identifiers) + nsID = ns.ID + } + + // Module + mod := n.relMod + if mod.ID <= 0 { + m := findComposeModuleR(state.ParentResources, n.res.ModRef.Identifiers) + mod.ID = m.ID + } + + // Some pointing + rm := n.res.RecMap + im := n.res.IDMap + + return n.res.Walker(func(r *resource.ComposeRecordRaw) error { + rec := &types.Record{ + ID: im[r.ID], + NamespaceID: nsID, + ModuleID: mod.ID, + CreatedAt: time.Now(), + } + exists := rm[r.ID] != nil + + if rec.ID <= 0 && exists { + rec.ID = rm[r.ID].ID + } else { + rec.ID = nextID() + } + + im[r.ID] = rec.ID + + if state.Conflicting { + return nil + } + + // Sys values + // + // @todo + for k, v := range r.SysValues { + if v == "" { + continue + } + + switch k { + case "createdAt": + // @todo set time + rec.CreatedAt = time.Now() + + case "updatedAt": + // @todo set time + rec.UpdatedAt = nil + + case "deletedAt": + // @todo set time + rec.DeletedAt = nil + } + } + + rvs := make(types.RecordValueSet, 0, len(r.Values)) + for k, v := range r.Values { + rv := &types.RecordValue{ + RecordID: rec.ID, + Name: k, + Value: v, + Updated: true, + } + + rvs = append(rvs, rv) + } + + // @todo validator + rec.Values = rvSanitizer.Run(mod, rvs) + + // Create a new record + if !exists { + err = store.CreateComposeRecord(ctx, s, mod, rec) + if err != nil { + return err + } + return nil + } + + // Update existing + err = store.UpdateComposeRecord(ctx, s, mod, rec) + if err != nil { + return err + } + + return nil + }) +} diff --git a/pkg/envoy/store/encoder.go b/pkg/envoy/store/encoder.go new file mode 100644 index 000000000..4742aad06 --- /dev/null +++ b/pkg/envoy/store/encoder.go @@ -0,0 +1,149 @@ +package store + +import ( + "context" + "errors" + "sync" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" + "github.com/davecgh/go-spew/spew" +) + +const ( + // Skip skips the existing resource + Skip mergeAlg = iota + // Replace replaces the existing resource + Replace + // MergeLeft updates the existing resource, giving priority to the existing data + MergeLeft + // MergeRight updates the existing resource, giving priority to the new data + MergeRight +) + +type ( + storeEncoder struct { + s store.Storer + cfg *EncoderConfig + + // Each resource should define its own state that is used when encoding the resource. + // Such approach removes the need for a janky generic global state. + // This also simplifies any slight deviations between resource complexities. + state map[resource.Interface]resourceState + } + + mergeAlg int + + // EncoderConfig allows us to configure the resource encoding process + EncoderConfig struct { + OnExisting mergeAlg + } + + // resourceState allows each conforming struct to be initialized and encoded + // by the store encoder + resourceState interface { + Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) + Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) + } +) + +// NewStoreEncoder initializes a fresh store encoder +// +// If no config is provided, it uses Skip as the default merge alg. +func NewStoreEncoder(s store.Storer, cfg *EncoderConfig) envoy.PrepareEncoder { + if cfg == nil { + cfg = &EncoderConfig{ + OnExisting: Skip, + } + } + + return &storeEncoder{ + s: s, + cfg: cfg, + + state: make(map[resource.Interface]resourceState), + } +} + +// Prepare prepares the encoder for the given set of resources +// +// It initializes and prepares the resource state for each provided resource +func (se *storeEncoder) Prepare(ctx context.Context, ee ...*envoy.ResourceState) (err error) { + f := func(rs resourceState, es *envoy.ResourceState) error { + err = rs.Prepare(ctx, se.s, es) + if err != nil { + return err + } + + se.state[es.Res] = rs + return nil + } + + for _, e := range ee { + switch res := e.Res.(type) { + // Compose things + case *resource.ComposeNamespace: + err = f(NewComposeNamespaceState(res, se.cfg), e) + case *resource.ComposeModule: + err = f(NewComposeModuleState(res, se.cfg), e) + case *resource.ComposeRecord: + err = f(NewComposeRecordState(res, se.cfg), e) + case *resource.ComposeChart: + err = f(NewComposeChartState(res, se.cfg), e) + case *resource.ComposePage: + err = f(NewComposePageState(res, se.cfg), e) + + // System things + case *resource.User: + err = f(NewUserState(res, se.cfg), e) + case *resource.Role: + err = f(NewRole(res, se.cfg), e) + case *resource.Application: + err = f(NewApplicationState(res, se.cfg), e) + case *resource.Settings: + err = f(NewSettingsState(res, se.cfg), e) + case *resource.RbacRule: + err = f(NewRbacRuleState(res, se.cfg), e) + + default: + return errors.New("[encoder] unknown resource; @todo error") + } + + if err != nil { + return err + } + } + + return nil +} + +// Encode encodes available resource states using the given store encoder +func (se *storeEncoder) Encode(ctx context.Context, wg *sync.WaitGroup, rc envoy.Rc, ec envoy.Ec) { + defer wg.Done() + + var e *envoy.ResourceState + err := store.Tx(ctx, se.s, func(ctx context.Context, s store.Storer) (err error) { + for { + e = <-rc + if e == nil { + return nil + } + + state := se.state[e.Res] + if state == nil { + return errors.New("Resource state not defined; @todo error") + } + + err = state.Encode(ctx, se.s, e) + if err != nil { + return err + } + } + }) + + if err != nil { + // ec <- err + spew.Dump(err) + } +} diff --git a/pkg/envoy/store/rbac_rule.go b/pkg/envoy/store/rbac_rule.go new file mode 100644 index 000000000..1d0e4bd8c --- /dev/null +++ b/pkg/envoy/store/rbac_rule.go @@ -0,0 +1,131 @@ +package store + +import ( + "context" + "errors" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/pkg/rbac" + "github.com/cortezaproject/corteza-server/store" + "github.com/cortezaproject/corteza-server/system/types" +) + +type ( + rbacRuleState struct { + cfg *EncoderConfig + + res *resource.RbacRule + rule *rbac.Rule + + relRole *types.Role + relResource resource.Interface + } +) + +var ( + // Use this as a cache, so we don't have to constantly fetch them + rbacRules rbac.RuleSet = nil +) + +func NewRbacRuleState(res *resource.RbacRule, cfg *EncoderConfig) resourceState { + return &rbacRuleState{ + cfg: cfg, + + res: res, + } +} + +func (n *rbacRuleState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Preload rbac rules if needed + if rbacRules == nil { + rbacRules, _, err = store.SearchRbacRules(ctx, s, rbac.RuleFilter{}) + if err != store.ErrNotFound && err != nil { + return err + } + } + + // Check for role + n.relRole, err = findRoleRS(ctx, s, state.ParentResources, n.res.RefRole.Identifiers) + if err != nil { + return err + } + if n.relRole == nil { + return errors.New("[rbac rule] couldn't resolve role; @todo error") + } + + // Check for resource if there is any + if n.res.RefResource != nil { + n.relResource = n.findResourceR(ctx, state.ParentResources, n.res.RefResource.ResourceType, n.res.RefResource.Identifiers) + if n.relResource == nil { + // Try to find it in the store + // @todo... + } + } + + return nil +} + +func (n *rbacRuleState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + rule := n.res.Res + + // Role + rule.RoleID = n.relRole.ID + if rule.RoleID <= 0 { + rl := findRoleR(state.ParentResources, n.res.RefRole.Identifiers) + rule.RoleID = rl.ID + } + + // Resource? + if n.res.RefResource != nil { + if ir, ok := n.relResource.(resource.IdentifiableInterface); !ok { + return errors.New("[rbac rule] resource not identifiable; @todo error") + } else { + rule.Resource = rule.Resource.AppendID(ir.SysID()) + } + } else if rule.Resource.IsAppendable() { + rule.Resource = rule.Resource.AppendWildcard() + } + + // There isn't anything to merge really, so Skip & MergeLeft skip it; + // Replace & MergeRight replace it. + // + // Replacing is handled after + switch n.cfg.OnExisting { + case Skip, + MergeLeft: + return nil + } + + nrr := rbacRules.Merge(rule) + d, u := nrr.Dirty() + + err = store.DeleteRbacRule(ctx, s, d...) + if err != nil { + return + } + + err = store.UpsertRbacRule(ctx, s, u...) + if err != nil { + return + } + + nrr.Clear() + rbacRules = nrr + return nil +} + +func (n *rbacRuleState) findResourceR(ctx context.Context, rr resource.InterfaceSet, rt string, ii resource.Identifiers) (rtr resource.Interface) { + for _, r := range rr { + if r.ResourceType() != rt { + continue + } + + if !r.Identifiers().HasAny(ii) { + continue + } + + return r + } + return nil +} diff --git a/pkg/envoy/store/role.go b/pkg/envoy/store/role.go new file mode 100644 index 000000000..79ba392a6 --- /dev/null +++ b/pkg/envoy/store/role.go @@ -0,0 +1,166 @@ +package store + +import ( + "context" + "time" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" + "github.com/cortezaproject/corteza-server/system/types" +) + +type ( + roleState struct { + cfg *EncoderConfig + + res *resource.Role + rl *types.Role + } +) + +func NewRole(res *resource.Role, cfg *EncoderConfig) resourceState { + return &roleState{ + cfg: cfg, + + res: res, + } +} + +func (n *roleState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + n.rl, err = findRoleS(ctx, s, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.rl != nil { + n.res.Res.ID = n.rl.ID + } + return nil +} + +func (n *roleState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + rl := n.res.Res + exists := n.rl != nil && n.rl.ID > 0 + + // Determine the ID + if rl.ID <= 0 && exists { + rl.ID = n.rl.ID + } + if rl.ID <= 0 { + rl.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Create a fresh role + if !exists { + return store.CreateRole(ctx, s, rl) + } + + // Update existing roles + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + rl = mergeRole(n.rl, rl) + + case MergeRight: + rl = mergeRole(rl, n.rl) + } + + err = store.UpdateRole(ctx, s, rl) + if err != nil { + return err + } + + n.res.Res = rl + return nil +} + +// mergeRole merges b into a, prioritising a +func mergeRole(a, b *types.Role) *types.Role { + c := *a + + if c.Name == "" { + c.Name = b.Name + } + if c.Handle == "" { + c.Handle = b.Handle + } + + return &c +} + +// findRoleRS looks for the role in the resources & the store +// +// Provided resources are prioritized. +func findRoleRS(ctx context.Context, s store.Storer, rr resource.InterfaceSet, ii resource.Identifiers) (rl *types.Role, err error) { + rl = findRoleR(rr, ii) + if rl != nil { + return rl, nil + } + + return findRoleS(ctx, s, makeGenericFilter(ii)) +} + +// findRoleS looks for the role in the store +func findRoleS(ctx context.Context, s store.Storer, gf genericFilter) (rl *types.Role, err error) { + if gf.id > 0 { + rl, err = store.LookupRoleByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if rl != nil { + return + } + } + + if gf.handle != "" { + rl, err = store.LookupRoleByHandle(ctx, s, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if rl != nil { + return + } + } + + return nil, nil +} + +// findRoleR looks for the role in the resources +func findRoleR(rr resource.InterfaceSet, ii resource.Identifiers) (rl *types.Role) { + var rlRes *resource.Role + var ok bool + + rr.Walk(func(r resource.Interface) error { + rlRes, ok = r.(*resource.Role) + if !ok { + return nil + } + + if rlRes.Identifiers().HasAny(r.Identifiers()) { + rlRes = nil + } + return nil + }) + + // Found it + if rlRes != nil { + return rlRes.Res + } + + return nil +} diff --git a/pkg/envoy/store/settings.go b/pkg/envoy/store/settings.go new file mode 100644 index 000000000..3b14e5f1e --- /dev/null +++ b/pkg/envoy/store/settings.go @@ -0,0 +1,90 @@ +package store + +import ( + "context" + "time" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" + "github.com/cortezaproject/corteza-server/system/types" +) + +type ( + settingsState struct { + cfg *EncoderConfig + + res *resource.Settings + ss types.SettingValueSet + } +) + +func NewSettingsState(res *resource.Settings, cfg *EncoderConfig) resourceState { + return &settingsState{ + cfg: cfg, + + res: res, + } +} + +func (n *settingsState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Preload settings + n.ss, _, err = store.SearchSettings(ctx, s, types.SettingsFilter{}) + if err == store.ErrNotFound { + n.ss = make(types.SettingValueSet, 0, len(n.res.Res)) + } else if err != nil { + return err + } + + // Default values + for _, s := range n.res.Res { + if s.UpdatedAt.IsZero() { + s.UpdatedAt = time.Now() + } + } + + // Nothing else to do. + // Settings can't conflict either. + + return nil +} + +func (n *settingsState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + ss := make(types.SettingValueSet, 0, len(n.res.Res)) + + for _, ns := range n.res.Res { + os := n.ss.First(ns.Name) + if os != nil { + // Update existing setting + switch n.cfg.OnExisting { + case Skip: + ss = append(ss, os) + + case Replace: + ss = append(ss, ns) + + case MergeLeft: + ss = append(ss, mergeSettings(os, ns)) + + case MergeRight: + ss = append(ss, mergeSettings(ns, os)) + } + } else { + // Create fresh setting + ss = append(ss, ns) + } + } + + return store.UpsertSetting(ctx, s, ss...) +} + +// mergeSettings merges b into a, prioritising a +func mergeSettings(a, b *types.SettingValue) *types.SettingValue { + c := *a + + if len(c.Value) <= 0 { + c.Value = b.Value + } + + return &c +} diff --git a/pkg/envoy/store/user.go b/pkg/envoy/store/user.go new file mode 100644 index 000000000..c20472f51 --- /dev/null +++ b/pkg/envoy/store/user.go @@ -0,0 +1,181 @@ +package store + +import ( + "context" + "time" + + "github.com/cortezaproject/corteza-server/pkg/envoy" + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/store" + "github.com/cortezaproject/corteza-server/system/types" +) + +type ( + userState struct { + cfg *EncoderConfig + + res *resource.User + u *types.User + } +) + +func NewUserState(res *resource.User, cfg *EncoderConfig) resourceState { + return &userState{ + cfg: cfg, + + res: res, + } +} + +func (n *userState) Prepare(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + // Initial values + if n.res.Res.CreatedAt.IsZero() { + n.res.Res.CreatedAt = time.Now() + } + + // Try to get the original user + // @todo make filtering more flexible (email, username, ...) + n.u, err = findUserS(ctx, s, makeGenericFilter(n.res.Identifiers())) + if err != nil { + return err + } + + if n.u != nil { + n.res.Res.ID = n.u.ID + } + return nil +} + +func (n *userState) Encode(ctx context.Context, s store.Storer, state *envoy.ResourceState) (err error) { + res := n.res.Res + exists := n.u != nil && n.u.ID > 0 + + // Determine the ID + if res.ID <= 0 && exists { + res.ID = n.u.ID + } + if res.ID <= 0 { + res.ID = nextID() + } + + // This is not possible, but let's do it anyway + if state.Conflicting { + return nil + } + + // Create a fresh user + if !exists { + return store.CreateUser(ctx, s, res) + } + + // Update existing user + switch n.cfg.OnExisting { + case Skip: + return nil + + case MergeLeft: + res = mergeUsers(n.u, res) + + case MergeRight: + res = mergeUsers(res, n.u) + } + + err = store.UpdateUser(ctx, s, res) + if err != nil { + return err + } + + n.res.Res = res + return nil +} + +// mergeUsers merges b into a, prioritising a +func mergeUsers(a, b *types.User) *types.User { + c := *a + + if c.Username == "" { + c.Username = b.Username + } + if c.Email == "" { + c.Email = b.Email + } + if c.Name == "" { + c.Name = b.Name + } + if c.Handle == "" { + c.Handle = b.Handle + } + if c.Kind == "" { + c.Kind = b.Kind + } + if *c.Meta == (types.UserMeta{}) { + c.Meta = b.Meta + } + + return &c +} + +// findUserRS looks for the user in the resources & the store +// +// Provided resources are prioritized. +func findUserRS(ctx context.Context, s store.Storer, rr resource.InterfaceSet, ii resource.Identifiers) (u *types.User, err error) { + u = findUserR(ctx, rr, ii) + if u != nil { + return u, nil + } + + return findUserS(ctx, s, makeGenericFilter(ii)) +} + +// findUserS looks for the user in the store +func findUserS(ctx context.Context, s store.Storer, gf genericFilter) (u *types.User, err error) { + if gf.id > 0 { + u, err = store.LookupUserByID(ctx, s, gf.id) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if u != nil { + return + } + } + + if gf.handle != "" { + u, err = store.LookupUserByHandle(ctx, s, gf.handle) + if err != nil && err != store.ErrNotFound { + return nil, err + } + + if u != nil { + return + } + } + + return nil, nil +} + +// findUserR looks for the user in the resources +func findUserR(ctx context.Context, rr resource.InterfaceSet, ii resource.Identifiers) (u *types.User) { + // Try to find it in the parent resources + var uRes *resource.User + var ok bool + + rr.Walk(func(r resource.Interface) error { + uRes, ok = r.(*resource.User) + if !ok { + return nil + } + + if !uRes.Identifiers().HasAny(r.Identifiers()) { + uRes = nil + } + return nil + }) + + // Found it + if uRes != nil { + return uRes.Res + } + + return nil +} diff --git a/pkg/envoy/store/util.go b/pkg/envoy/store/util.go new file mode 100644 index 000000000..051b7aff3 --- /dev/null +++ b/pkg/envoy/store/util.go @@ -0,0 +1,54 @@ +package store + +import ( + "regexp" + "strconv" + + "github.com/cortezaproject/corteza-server/pkg/envoy/resource" + "github.com/cortezaproject/corteza-server/pkg/handle" + "github.com/cortezaproject/corteza-server/pkg/id" +) + +type ( + genericFilter struct { + id uint64 + handle string + name string + } +) + +var ( + // simple uint check. + // we'll use the pkg/handle to check for handles. + refy = regexp.MustCompile(`^[1-9](\d*)$`) + + // wrapper around nextID that will aid service testing + nextID = func() uint64 { + return id.Next() + } +) + +// makeGenericFilter is a helper to determine the base resource filter. +// +// It attempts to determine an identifier, handle, and name. +func makeGenericFilter(ii resource.Identifiers) (f genericFilter) { + for i := range ii { + if i == "" { + continue + } + + if refy.MatchString(i) && f.id <= 0 { + id, err := strconv.ParseUint(i, 10, 64) + if err != nil { + continue + } + f.id = id + } else if handle.IsValid(i) && f.handle == "" { + f.handle = i + } else if f.name == "" { + f.name = i + } + } + + return f +}