diff --git a/federation/service/node_sync.go b/federation/service/node_sync.go new file mode 100644 index 000000000..9415b5600 --- /dev/null +++ b/federation/service/node_sync.go @@ -0,0 +1,66 @@ +package service + +import ( + "context" + + "github.com/cortezaproject/corteza-server/federation/types" + "github.com/cortezaproject/corteza-server/pkg/actionlog" + "github.com/cortezaproject/corteza-server/pkg/filter" + "github.com/cortezaproject/corteza-server/store" +) + +type ( + nodeSync struct { + store store.Storer + actionlog actionlog.Recorder + } + + NodeSyncService interface { + Create(ctx context.Context, new *types.NodeSync) (*types.NodeSync, error) + LookupLastSuccessfulSync(ctx context.Context, nodeID uint64, syncType string) (*types.NodeSync, error) + } +) + +func NodeSync() NodeSyncService { + return &nodeSync{ + store: DefaultStore, + actionlog: DefaultActionlog, + } +} + +func (svc nodeSync) Create(ctx context.Context, new *types.NodeSync) (*types.NodeSync, error) { + var ( + aProps = &nodeSyncActionProps{nodeSync: new} + ) + + err := store.Tx(ctx, svc.store, func(ctx context.Context, s store.Storer) (err error) { + if _, err := DefaultNode.FindByID(ctx, new.NodeID); err != nil { + return NodeSyncErrNodeNotFound() + } + + return store.CreateFederationNodesSync(ctx, s, new) + }) + + return new, svc.recordAction(ctx, aProps, NodeSyncActionCreate, err) +} + +func (svc nodeSync) LookupLastSuccessfulSync(ctx context.Context, nodeID uint64, syncType string) (ns *types.NodeSync, err error) { + // todo - filter by synctype does not work + s, _, err := store.SearchFederationNodesSyncs(ctx, svc.store, types.NodeSyncFilter{ + NodeID: nodeID, + SyncType: syncType, + SyncStatus: types.NodeSyncStatusSuccess, + Sorting: filter.Sorting{ + Sort: filter.SortExprSet{ + &filter.SortExpr{Column: "time_action", Descending: true}, + }, + }, + Paging: filter.Paging{Limit: 1}, + }) + + if err != nil || len(s) == 0 { + return nil, err + } + + return s[0], nil +} diff --git a/federation/service/node_sync_actions.gen.go b/federation/service/node_sync_actions.gen.go new file mode 100644 index 000000000..360b420aa --- /dev/null +++ b/federation/service/node_sync_actions.gen.go @@ -0,0 +1,527 @@ +package service + +// This file is auto-generated. +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. +// +// Definitions file that controls how this file is generated: +// federation/service/node_sync_actions.yaml + +import ( + "context" + "errors" + "fmt" + "strings" + "time" + + "github.com/cortezaproject/corteza-server/federation/types" + "github.com/cortezaproject/corteza-server/pkg/actionlog" +) + +type ( + nodeSyncActionProps struct { + nodeSync *types.NodeSync + nodeSyncFilter *types.NodeSyncFilter + } + + nodeSyncAction struct { + timestamp time.Time + resource string + action string + log string + severity actionlog.Severity + + // prefix for error when action fails + errorMessage string + + props *nodeSyncActionProps + } + + nodeSyncError struct { + timestamp time.Time + error string + resource string + action string + message string + log string + severity actionlog.Severity + + wrap error + + props *nodeSyncActionProps + } +) + +var ( + // just a placeholder to cover template cases w/o fmt package use + _ = fmt.Println +) + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Props methods +// setNodeSync updates nodeSyncActionProps's nodeSync +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *nodeSyncActionProps) setNodeSync(nodeSync *types.NodeSync) *nodeSyncActionProps { + p.nodeSync = nodeSync + return p +} + +// setNodeSyncFilter updates nodeSyncActionProps's nodeSyncFilter +// +// Allows method chaining +// +// This function is auto-generated. +// +func (p *nodeSyncActionProps) setNodeSyncFilter(nodeSyncFilter *types.NodeSyncFilter) *nodeSyncActionProps { + p.nodeSyncFilter = nodeSyncFilter + return p +} + +// serialize converts nodeSyncActionProps to actionlog.Meta +// +// This function is auto-generated. +// +func (p nodeSyncActionProps) serialize() actionlog.Meta { + var ( + m = make(actionlog.Meta) + ) + + if p.nodeSync != nil { + m.Set("nodeSync.NodeID", p.nodeSync.NodeID, true) + m.Set("nodeSync.SyncStatus", p.nodeSync.SyncStatus, true) + m.Set("nodeSync.SyncType", p.nodeSync.SyncType, true) + m.Set("nodeSync.TimeOfAction", p.nodeSync.TimeOfAction, true) + } + if p.nodeSyncFilter != nil { + m.Set("nodeSyncFilter.Query", p.nodeSyncFilter.Query, true) + } + + return m +} + +// tr translates string and replaces meta value placeholder with values +// +// This function is auto-generated. +// +func (p nodeSyncActionProps) tr(in string, err error) string { + var ( + pairs = []string{"{err}"} + // first non-empty string + fns = func(ii ...interface{}) string { + for _, i := range ii { + if s := fmt.Sprintf("%v", i); len(s) > 0 { + return s + } + } + + return "" + } + ) + + if err != nil { + for { + // Unwrap errors + ue := errors.Unwrap(err) + if ue == nil { + break + } + + err = ue + } + + pairs = append(pairs, err.Error()) + } else { + pairs = append(pairs, "nil") + } + + if p.nodeSync != nil { + // replacement for "{nodeSync}" (in order how fields are defined) + pairs = append( + pairs, + "{nodeSync}", + fns( + p.nodeSync.NodeID, + p.nodeSync.SyncStatus, + p.nodeSync.SyncType, + p.nodeSync.TimeOfAction, + ), + ) + pairs = append(pairs, "{nodeSync.NodeID}", fns(p.nodeSync.NodeID)) + pairs = append(pairs, "{nodeSync.SyncStatus}", fns(p.nodeSync.SyncStatus)) + pairs = append(pairs, "{nodeSync.SyncType}", fns(p.nodeSync.SyncType)) + pairs = append(pairs, "{nodeSync.TimeOfAction}", fns(p.nodeSync.TimeOfAction)) + } + + if p.nodeSyncFilter != nil { + // replacement for "{nodeSyncFilter}" (in order how fields are defined) + pairs = append( + pairs, + "{nodeSyncFilter}", + fns( + p.nodeSyncFilter.Query, + ), + ) + pairs = append(pairs, "{nodeSyncFilter.Query}", fns(p.nodeSyncFilter.Query)) + } + return strings.NewReplacer(pairs...).Replace(in) +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Action methods + +// String returns loggable description as string +// +// This function is auto-generated. +// +func (a *nodeSyncAction) String() string { + var props = &nodeSyncActionProps{} + + if a.props != nil { + props = a.props + } + + return props.tr(a.log, nil) +} + +func (e *nodeSyncAction) LoggableAction() *actionlog.Action { + return &actionlog.Action{ + Timestamp: e.timestamp, + Resource: e.resource, + Action: e.action, + Severity: e.severity, + Description: e.String(), + Meta: e.props.serialize(), + } +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Error methods + +// String returns loggable description as string +// +// It falls back to message if log is not set +// +// This function is auto-generated. +// +func (e *nodeSyncError) String() string { + var props = &nodeSyncActionProps{} + + if e.props != nil { + props = e.props + } + + if e.wrap != nil && !strings.Contains(e.log, "{err}") { + // Suffix error log with {err} to ensure + // we log the cause for this error + e.log += ": {err}" + } + + return props.tr(e.log, e.wrap) +} + +// Error satisfies +// +// This function is auto-generated. +// +func (e *nodeSyncError) Error() string { + var props = &nodeSyncActionProps{} + + if e.props != nil { + props = e.props + } + + return props.tr(e.message, e.wrap) +} + +// Is fn for error equality check +// +// This function is auto-generated. +// +func (e *nodeSyncError) Is(err error) bool { + t, ok := err.(*nodeSyncError) + if !ok { + return false + } + + return t.resource == e.resource && t.error == e.error +} + +// Is fn for error equality check +// +// This function is auto-generated. +// +func (e *nodeSyncError) IsGeneric() bool { + return e.error == "generic" +} + +// Wrap wraps nodeSyncError around another error +// +// This function is auto-generated. +// +func (e *nodeSyncError) Wrap(err error) *nodeSyncError { + e.wrap = err + return e +} + +// Unwrap returns wrapped error +// +// This function is auto-generated. +// +func (e *nodeSyncError) Unwrap() error { + return e.wrap +} + +func (e *nodeSyncError) LoggableAction() *actionlog.Action { + return &actionlog.Action{ + Timestamp: e.timestamp, + Resource: e.resource, + Action: e.action, + Severity: e.severity, + Description: e.String(), + Error: e.Error(), + Meta: e.props.serialize(), + } +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Action constructors + +// NodeSyncActionLookup returns "federation:node_sync.lookup" error +// +// This function is auto-generated. +// +func NodeSyncActionLookup(props ...*nodeSyncActionProps) *nodeSyncAction { + a := &nodeSyncAction{ + timestamp: time.Now(), + resource: "federation:node_sync", + action: "lookup", + log: "looked-up for the last successful sync", + severity: actionlog.Info, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// NodeSyncActionCreate returns "federation:node_sync.create" error +// +// This function is auto-generated. +// +func NodeSyncActionCreate(props ...*nodeSyncActionProps) *nodeSyncAction { + a := &nodeSyncAction{ + timestamp: time.Now(), + resource: "federation:node_sync", + action: "create", + log: "created node_sync", + severity: actionlog.Notice, + } + + if len(props) > 0 { + a.props = props[0] + } + + return a +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* +// Error constructors + +// NodeSyncErrGeneric returns "federation:node_sync.generic" audit event as actionlog.Error +// +// +// This function is auto-generated. +// +func NodeSyncErrGeneric(props ...*nodeSyncActionProps) *nodeSyncError { + var e = &nodeSyncError{ + timestamp: time.Now(), + resource: "federation:node_sync", + error: "generic", + action: "error", + message: "failed to complete request due to internal error", + log: "{err}", + severity: actionlog.Error, + props: func() *nodeSyncActionProps { + if len(props) > 0 { + return props[0] + } + return nil + }(), + } + + if len(props) > 0 { + e.props = props[0] + } + + return e + +} + +// NodeSyncErrNotFound returns "federation:node_sync.notFound" audit event as actionlog.Warning +// +// +// This function is auto-generated. +// +func NodeSyncErrNotFound(props ...*nodeSyncActionProps) *nodeSyncError { + var e = &nodeSyncError{ + timestamp: time.Now(), + resource: "federation:node_sync", + error: "notFound", + action: "error", + message: "node_sync does not exist", + log: "node_sync does not exist", + severity: actionlog.Warning, + props: func() *nodeSyncActionProps { + if len(props) > 0 { + return props[0] + } + return nil + }(), + } + + if len(props) > 0 { + e.props = props[0] + } + + return e + +} + +// NodeSyncErrNodeNotFound returns "federation:node_sync.nodeNotFound" audit event as actionlog.Warning +// +// +// This function is auto-generated. +// +func NodeSyncErrNodeNotFound(props ...*nodeSyncActionProps) *nodeSyncError { + var e = &nodeSyncError{ + timestamp: time.Now(), + resource: "federation:node_sync", + error: "nodeNotFound", + action: "error", + message: "node does not exist", + log: "node does not exist", + severity: actionlog.Warning, + props: func() *nodeSyncActionProps { + if len(props) > 0 { + return props[0] + } + return nil + }(), + } + + if len(props) > 0 { + e.props = props[0] + } + + return e + +} + +// ********************************************************************************************************************* +// ********************************************************************************************************************* + +// recordAction is a service helper function wraps function that can return error +// +// context is used to enrich audit log entry with current user info, request ID, IP address... +// props are collected action/error properties +// action (optional) fn will be used to construct nodeSyncAction struct from given props (and error) +// err is any error that occurred while action was happening +// +// Action has success and fail (error) state: +// - when recorded without an error (4th param), action is recorded as successful. +// - when an additional error is given (4th param), action is used to wrap +// the additional error +// +// This function is auto-generated. +// +func (svc nodeSync) recordAction(ctx context.Context, props *nodeSyncActionProps, action func(...*nodeSyncActionProps) *nodeSyncAction, err error) error { + var ( + ok bool + + // Return error + retError *nodeSyncError + + // Recorder error + recError *nodeSyncError + ) + + if err != nil { + if retError, ok = err.(*nodeSyncError); !ok { + // got non-nodeSync error, wrap it with NodeSyncErrGeneric + retError = NodeSyncErrGeneric(props).Wrap(err) + + if action != nil { + // copy action to returning and recording error + retError.action = action().action + } + + // we'll use NodeSyncErrGeneric for recording too + // because it can hold more info + recError = retError + } else if retError != nil { + if action != nil { + // copy action to returning and recording error + retError.action = action().action + } + // start with copy of return error for recording + // this will be updated with tha root cause as we try and + // unwrap the error + recError = retError + + // find the original recError for this error + // for the purpose of logging + var unwrappedError error = retError + for { + if unwrappedError = errors.Unwrap(unwrappedError); unwrappedError == nil { + // nothing wrapped + break + } + + // update recError ONLY of wrapped error is of type nodeSyncError + if unwrappedSinkError, ok := unwrappedError.(*nodeSyncError); ok { + recError = unwrappedSinkError + } + } + + if retError.props == nil { + // set props on returning error if empty + retError.props = props + } + + if recError.props == nil { + // set props on recording error if empty + recError.props = props + } + } + } + + if svc.actionlog != nil { + if retError != nil { + // failed action, log error + svc.actionlog.Record(ctx, recError) + } else if action != nil { + // successful + svc.actionlog.Record(ctx, action(props)) + } + } + + if err == nil { + // retError not an interface and that WILL (!!) cause issues + // with nil check (== nil) when it is not explicitly returned + return nil + } + + return retError +} diff --git a/federation/service/node_sync_actions.yaml b/federation/service/node_sync_actions.yaml new file mode 100644 index 000000000..503bdd147 --- /dev/null +++ b/federation/service/node_sync_actions.yaml @@ -0,0 +1,38 @@ +# List of loggable service actions + +resource: federation:node_sync +service: nodeSync + +# Default sensitivity for actions +defaultActionSeverity: notice + +# default severity for errors +defaultErrorSeverity: error + +import: + - github.com/cortezaproject/corteza-server/federation/types + +props: + - name: nodeSync + type: "*types.NodeSync" + fields: [ NodeID, SyncStatus, SyncType, TimeOfAction ] + - name: nodeSyncFilter + type: "*types.NodeSyncFilter" + fields: [ Query ] + +actions: + - action: lookup + log: "looked-up for the last successful sync" + severity: info + + - action: create + log: "created node_sync" + +errors: + - error: notFound + message: "node_sync does not exist" + severity: warning + + - error: nodeNotFound + message: "node does not exist" + severity: warning diff --git a/federation/service/service.go b/federation/service/service.go index 1a32ff8e7..5f33dbb75 100644 --- a/federation/service/service.go +++ b/federation/service/service.go @@ -41,6 +41,7 @@ var ( DefaultActionlog actionlog.Recorder DefaultNode *node + DefaultNodeSync NodeSyncService DefaultExposedModule ExposedModuleService DefaultSharedModule SharedModuleService @@ -119,6 +120,7 @@ func Initialize(ctx context.Context, log *zap.Logger, s store.Storer, c Config) hcd.Add(objstore.Healthcheck(DefaultObjectStore), "Store/Federation") DefaultNode = Node(DefaultStore, service.DefaultUser, DefaultActionlog, auth.DefaultJwtHandler) + DefaultNodeSync = NodeSync() DefaultExposedModule = ExposedModule() DefaultSharedModule = SharedModule() diff --git a/federation/types/node_sync.go b/federation/types/node_sync.go new file mode 100644 index 000000000..286eca772 --- /dev/null +++ b/federation/types/node_sync.go @@ -0,0 +1,37 @@ +package types + +import ( + "time" + + "github.com/cortezaproject/corteza-server/pkg/filter" +) + +var ( + NodeSyncTypeStructure = "sync_structure" + NodeSyncTypeData = "sync_data" + NodeSyncStatusSuccess = "success" + NodeSyncStatusError = "error" +) + +type ( + NodeSync struct { + NodeID uint64 `json:"nodeID,string"` + SyncStatus string `json:"syncStatus"` + SyncType string `json:"syncType"` + + TimeOfAction time.Time `json:"timeOfAction"` + } + + NodeSyncFilter struct { + NodeID uint64 `json:"nodeID"` + SyncStatus string `json:"syncStatus"` + SyncType string `json:"syncType"` + + Query string `json:"name"` + + Check func(*NodeSync) (bool, error) `json:"-"` + + filter.Sorting + filter.Paging + } +) diff --git a/federation/types/type_set.gen.go b/federation/types/type_set.gen.go index 1dc9ef0d0..871626cb2 100644 --- a/federation/types/type_set.gen.go +++ b/federation/types/type_set.gen.go @@ -25,6 +25,11 @@ type ( // This type is auto-generated. NodeSet []*Node + // NodeSyncSet slice of NodeSync + // + // This type is auto-generated. + NodeSyncSet []*NodeSync + // SharedModuleSet slice of SharedModule // // This type is auto-generated. @@ -173,6 +178,36 @@ func (set NodeSet) IDs() (IDs []uint64) { return } +// Walk iterates through every slice item and calls w(NodeSync) err +// +// This function is auto-generated. +func (set NodeSyncSet) Walk(w func(*NodeSync) error) (err error) { + for i := range set { + if err = w(set[i]); err != nil { + return + } + } + + return +} + +// Filter iterates through every slice item, calls f(NodeSync) (bool, err) and return filtered slice +// +// This function is auto-generated. +func (set NodeSyncSet) Filter(f func(*NodeSync) (bool, error)) (out NodeSyncSet, err error) { + var ok bool + out = NodeSyncSet{} + for i := range set { + if ok, err = f(set[i]); err != nil { + return + } else if ok { + out = append(out, set[i]) + } + } + + return +} + // Walk iterates through every slice item and calls w(SharedModule) err // // This function is auto-generated. diff --git a/federation/types/type_set.gen_test.go b/federation/types/type_set.gen_test.go index e1767861b..7bb5af4bf 100644 --- a/federation/types/type_set.gen_test.go +++ b/federation/types/type_set.gen_test.go @@ -250,6 +250,62 @@ func TestNodeSetIDs(t *testing.T) { } } +func TestNodeSyncSetWalk(t *testing.T) { + var ( + value = make(NodeSyncSet, 3) + req = require.New(t) + ) + + // check walk with no errors + { + err := value.Walk(func(*NodeSync) error { + return nil + }) + req.NoError(err) + } + + // check walk with error + req.Error(value.Walk(func(*NodeSync) error { return fmt.Errorf("walk error") })) +} + +func TestNodeSyncSetFilter(t *testing.T) { + var ( + value = make(NodeSyncSet, 3) + req = require.New(t) + ) + + // filter nothing + { + set, err := value.Filter(func(*NodeSync) (bool, error) { + return true, nil + }) + req.NoError(err) + req.Equal(len(set), len(value)) + } + + // filter one item + { + found := false + set, err := value.Filter(func(*NodeSync) (bool, error) { + if !found { + found = true + return found, nil + } + return false, nil + }) + req.NoError(err) + req.Len(set, 1) + } + + // filter error + { + _, err := value.Filter(func(*NodeSync) (bool, error) { + return false, fmt.Errorf("filter error") + }) + req.Error(err) + } +} + func TestSharedModuleSetWalk(t *testing.T) { var ( value = make(SharedModuleSet, 3) diff --git a/federation/types/types.yaml b/federation/types/types.yaml index 6ff7de46a..1ea5bb408 100644 --- a/federation/types/types.yaml +++ b/federation/types/types.yaml @@ -1,5 +1,7 @@ types: Node: + NodeSync: + noIdField: true ExposedModule: SharedModule: ModuleMapping: diff --git a/store/federation_nodes_sync.gen.go b/store/federation_nodes_sync.gen.go new file mode 100644 index 000000000..0981b1582 --- /dev/null +++ b/store/federation_nodes_sync.gen.go @@ -0,0 +1,85 @@ +package store + +// This file is auto-generated. +// +// Template: pkg/codegen/assets/store_base.gen.go.tpl +// Definitions: store/federation_nodes_sync.yaml +// +// Changes to this file may cause incorrect behavior and will be lost if +// the code is regenerated. + +import ( + "context" + "github.com/cortezaproject/corteza-server/federation/types" +) + +type ( + FederationNodesSyncs interface { + SearchFederationNodesSyncs(ctx context.Context, f types.NodeSyncFilter) (types.NodeSyncSet, types.NodeSyncFilter, error) + LookupFederationNodesSyncByNodeID(ctx context.Context, node_id uint64) (*types.NodeSync, error) + LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus(ctx context.Context, node_id uint64, sync_type string, sync_status string) (*types.NodeSync, error) + + CreateFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) error + + UpdateFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) error + + UpsertFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) error + + DeleteFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) error + DeleteFederationNodesSyncByNodeID(ctx context.Context, nodeID uint64) error + + TruncateFederationNodesSyncs(ctx context.Context) error + } +) + +var _ *types.NodeSync +var _ context.Context + +// SearchFederationNodesSyncs returns all matching FederationNodesSyncs from store +func SearchFederationNodesSyncs(ctx context.Context, s FederationNodesSyncs, f types.NodeSyncFilter) (types.NodeSyncSet, types.NodeSyncFilter, error) { + return s.SearchFederationNodesSyncs(ctx, f) +} + +// LookupFederationNodesSyncByNodeID searches for sync activity by node ID +// +// It returns sync activity +func LookupFederationNodesSyncByNodeID(ctx context.Context, s FederationNodesSyncs, node_id uint64) (*types.NodeSync, error) { + return s.LookupFederationNodesSyncByNodeID(ctx, node_id) +} + +// LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus searches for activity by node, type and status +// +// It returns sync activity +func LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus(ctx context.Context, s FederationNodesSyncs, node_id uint64, sync_type string, sync_status string) (*types.NodeSync, error) { + return s.LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus(ctx, node_id, sync_type, sync_status) +} + +// CreateFederationNodesSync creates one or more FederationNodesSyncs in store +func CreateFederationNodesSync(ctx context.Context, s FederationNodesSyncs, rr ...*types.NodeSync) error { + return s.CreateFederationNodesSync(ctx, rr...) +} + +// UpdateFederationNodesSync updates one or more (existing) FederationNodesSyncs in store +func UpdateFederationNodesSync(ctx context.Context, s FederationNodesSyncs, rr ...*types.NodeSync) error { + return s.UpdateFederationNodesSync(ctx, rr...) +} + +// UpsertFederationNodesSync creates new or updates existing one or more FederationNodesSyncs in store +func UpsertFederationNodesSync(ctx context.Context, s FederationNodesSyncs, rr ...*types.NodeSync) error { + return s.UpsertFederationNodesSync(ctx, rr...) +} + +// DeleteFederationNodesSync Deletes one or more FederationNodesSyncs from store +func DeleteFederationNodesSync(ctx context.Context, s FederationNodesSyncs, rr ...*types.NodeSync) error { + return s.DeleteFederationNodesSync(ctx, rr...) +} + +// DeleteFederationNodesSyncByNodeID Deletes FederationNodesSync from store +func DeleteFederationNodesSyncByNodeID(ctx context.Context, s FederationNodesSyncs, nodeID uint64) error { + return s.DeleteFederationNodesSyncByNodeID(ctx, nodeID) +} + +// TruncateFederationNodesSyncs Deletes all FederationNodesSyncs from store +func TruncateFederationNodesSyncs(ctx context.Context, s FederationNodesSyncs) error { + return s.TruncateFederationNodesSyncs(ctx) +} diff --git a/store/federation_nodes_sync.yaml b/store/federation_nodes_sync.yaml new file mode 100644 index 000000000..2684b1e19 --- /dev/null +++ b/store/federation_nodes_sync.yaml @@ -0,0 +1,38 @@ +import: + - github.com/cortezaproject/corteza-server/federation/types + +types: + type: types.NodeSync + +fields: + - { field: NodeID, isPrimaryKey: true, sortable: true } + - { field: SyncType } + - { field: SyncStatus } + - { field: TimeOfAction, sortable: true } + +lookups: + - fields: [NodeID] + description: |- + searches for sync activity by node ID + + It returns sync activity + + - fields: [NodeID, SyncType, SyncStatus] + description: |- + searches for activity by node, type and status + + It returns sync activity + +search: + enablePaging: true + enableSorting: true + +rdbms: + alias: fdns + table: federation_nodes_sync + customFilterConverter: true + mapFields: + NodeID: { column: node_id } + SyncType: { column: sync_type } + SyncStatus: { column: sync_status } + TimeOfAction: { column: time_action } diff --git a/store/interfaces.gen.go b/store/interfaces.gen.go index 517c3b26d..e2ca065f6 100644 --- a/store/interfaces.gen.go +++ b/store/interfaces.gen.go @@ -19,6 +19,7 @@ package store // - store/federation_exposed_modules.yaml // - store/federation_module_mappings.yaml // - store/federation_nodes.yaml +// - store/federation_nodes_sync.yaml // - store/federation_shared_modules.yaml // - store/labels.yaml // - store/messaging_attachments.yaml @@ -58,6 +59,7 @@ type ( FederationExposedModules FederationModuleMappings FederationNodes + FederationNodesSyncs FederationSharedModules Labels MessagingAttachments diff --git a/store/rdbms/federation_nodes_sync.gen.go b/store/rdbms/federation_nodes_sync.gen.go new file mode 100644 index 000000000..83f403a31 --- /dev/null +++ b/store/rdbms/federation_nodes_sync.gen.go @@ -0,0 +1,523 @@ +package rdbms + +// This file is an auto-generated file +// +// Template: pkg/codegen/assets/store_rdbms.gen.go.tpl +// Definitions: store/federation_nodes_sync.yaml +// +// Changes to this file may cause incorrect behavior +// and will be lost if the code is regenerated. + +import ( + "context" + "database/sql" + "errors" + "fmt" + "github.com/Masterminds/squirrel" + "github.com/cortezaproject/corteza-server/federation/types" + "github.com/cortezaproject/corteza-server/pkg/filter" + "github.com/cortezaproject/corteza-server/store" +) + +var _ = errors.Is + +// SearchFederationNodesSyncs returns all matching rows +// +// This function calls convertFederationNodesSyncFilter with the given +// types.NodeSyncFilter and expects to receive a working squirrel.SelectBuilder +func (s Store) SearchFederationNodesSyncs(ctx context.Context, f types.NodeSyncFilter) (types.NodeSyncSet, types.NodeSyncFilter, error) { + var ( + err error + set []*types.NodeSync + q squirrel.SelectBuilder + ) + q, err = s.convertFederationNodesSyncFilter(f) + if err != nil { + return nil, f, err + } + + // Cleanup anything we've accidentally received... + f.PrevPage, f.NextPage = nil, nil + + // 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 + reversedCursor := f.PageCursor != nil && f.PageCursor.Reverse + + // If paging with reverse cursor, change the sorting + // direction for all columns we're sorting by + curSort := f.Sort.Clone() + if reversedCursor { + curSort.Reverse() + } + + return set, f, s.config.ErrorHandler(func() error { + set, err = s.fetchFullPageOfFederationNodesSyncs(ctx, q, curSort, f.PageCursor, f.Limit, f.Check) + + if err != nil { + return err + } + + if f.Limit > 0 && len(set) > 0 { + if f.PageCursor != nil && (!f.PageCursor.Reverse || uint(len(set)) == f.Limit) { + f.PrevPage = s.collectFederationNodesSyncCursorValues(set[0], curSort.Columns()...) + f.PrevPage.Reverse = true + } + + // Less items fetched then requested by page-limit + // not very likely there's another page + f.NextPage = s.collectFederationNodesSyncCursorValues(set[len(set)-1], curSort.Columns()...) + } + + f.PageCursor = nil + return nil + }()) +} + +// fetchFullPageOfFederationNodesSyncs collects all requested results. +// +// Function applies: +// - cursor conditions (where ...) +// - sorting rules (order by ...) +// - 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 +func (s Store) fetchFullPageOfFederationNodesSyncs( + ctx context.Context, + q squirrel.SelectBuilder, + sort filter.SortExprSet, + cursor *filter.PagingCursor, + limit uint, + check func(*types.NodeSync) (bool, error), +) ([]*types.NodeSync, error) { + var ( + set = make([]*types.NodeSync, 0, DefaultSliceCapacity) + aux []*types.NodeSync + last *types.NodeSync + + // 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 + reversedCursor = cursor != nil && cursor.Reverse + + // copy of the select builder + tryQuery squirrel.SelectBuilder + + fetched uint + err error + ) + + // Make sure we always end our sort by primary keys + if sort.Get("node_id") == nil { + sort = append(sort, &filter.SortExpr{Column: "node_id"}) + } + + // Apply sorting expr from filter to query + if q, err = setOrderBy(q, sort, s.sortableFederationNodesSyncColumns()...); err != nil { + return nil, err + } + + for try := 0; try < MaxRefetches; try++ { + tryQuery = setCursorCond(q, cursor) + if limit > 0 { + tryQuery = tryQuery.Limit(uint64(limit)) + } + + if aux, fetched, last, err = s.QueryFederationNodesSyncs(ctx, tryQuery, check); err != nil { + return nil, err + } + + if limit > 0 && uint(len(aux)) >= limit { + // we should use only as much as requested + set = append(set, aux[0:limit]...) + break + } else { + set = append(set, aux...) + } + + // if limit is not set or we've already collected enough items + // we can break the loop right away + if limit == 0 || fetched == 0 || fetched < limit { + break + } + + // 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 + if limit < MinEnsureFetchLimit { + limit = MinEnsureFetchLimit + } + + // @todo improve strategy for collecting next page with lower limit + + // Point cursor to the last fetched element + if cursor = s.collectFederationNodesSyncCursorValues(last, sort.Columns()...); cursor == nil { + break + } + } + + if reversedCursor { + // Cursor for previous page was used + // Fetched set needs to be reverseCursor because we've forced a descending order to + // get the previous page + for i, j := 0, len(set)-1; i < j; i, j = i+1, j-1 { + set[i], set[j] = set[j], set[i] + } + } + + return set, nil +} + +// QueryFederationNodesSyncs queries the database, converts and checks each row and +// returns collected set +// +// Fn also returns total number of fetched items and last fetched item so that the caller can construct cursor +// for next page of results +func (s Store) QueryFederationNodesSyncs( + ctx context.Context, + q squirrel.Sqlizer, + check func(*types.NodeSync) (bool, error), +) ([]*types.NodeSync, uint, *types.NodeSync, error) { + var ( + set = make([]*types.NodeSync, 0, DefaultSliceCapacity) + res *types.NodeSync + + // Query rows with + rows, err = s.Query(ctx, q) + + fetched uint + ) + + if err != nil { + return nil, 0, nil, err + } + + defer rows.Close() + for rows.Next() { + fetched++ + if err = rows.Err(); err == nil { + res, err = s.internalFederationNodesSyncRowScanner(rows) + } + + if err != nil { + return nil, 0, nil, err + } + + // If check function is set, call it and act accordingly + if check != nil { + if chk, err := check(res); err != nil { + return nil, 0, nil, err + } else if !chk { + // did not pass the check + // go with the next row + continue + } + } + + set = append(set, res) + } + + return set, fetched, res, rows.Err() +} + +// LookupFederationNodesSyncByNodeID searches for sync activity by node ID +// +// It returns sync activity +func (s Store) LookupFederationNodesSyncByNodeID(ctx context.Context, node_id uint64) (*types.NodeSync, error) { + return s.execLookupFederationNodesSync(ctx, squirrel.Eq{ + s.preprocessColumn("fdns.node_id", ""): store.PreprocessValue(node_id, ""), + }) +} + +// LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus searches for activity by node, type and status +// +// It returns sync activity +func (s Store) LookupFederationNodesSyncByNodeIDSyncTypeSyncStatus(ctx context.Context, node_id uint64, sync_type string, sync_status string) (*types.NodeSync, error) { + return s.execLookupFederationNodesSync(ctx, squirrel.Eq{ + s.preprocessColumn("fdns.node_id", ""): store.PreprocessValue(node_id, ""), + s.preprocessColumn("fdns.sync_type", ""): store.PreprocessValue(sync_type, ""), + s.preprocessColumn("fdns.sync_status", ""): store.PreprocessValue(sync_status, ""), + }) +} + +// CreateFederationNodesSync creates one or more rows in federation_nodes_sync table +func (s Store) CreateFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) (err error) { + for _, res := range rr { + err = s.checkFederationNodesSyncConstraints(ctx, res) + if err != nil { + return err + } + + err = s.execCreateFederationNodesSyncs(ctx, s.internalFederationNodesSyncEncoder(res)) + if err != nil { + return err + } + } + + return +} + +// UpdateFederationNodesSync updates one or more existing rows in federation_nodes_sync +func (s Store) UpdateFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) error { + return s.config.ErrorHandler(s.partialFederationNodesSyncUpdate(ctx, nil, rr...)) +} + +// partialFederationNodesSyncUpdate updates one or more existing rows in federation_nodes_sync +func (s Store) partialFederationNodesSyncUpdate(ctx context.Context, onlyColumns []string, rr ...*types.NodeSync) (err error) { + for _, res := range rr { + err = s.checkFederationNodesSyncConstraints(ctx, res) + if err != nil { + return err + } + + err = s.execUpdateFederationNodesSyncs( + ctx, + squirrel.Eq{ + s.preprocessColumn("fdns.node_id", ""): store.PreprocessValue(res.NodeID, ""), + }, + s.internalFederationNodesSyncEncoder(res).Skip("node_id").Only(onlyColumns...)) + if err != nil { + return s.config.ErrorHandler(err) + } + } + + return +} + +// UpsertFederationNodesSync updates one or more existing rows in federation_nodes_sync +func (s Store) UpsertFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) (err error) { + for _, res := range rr { + err = s.checkFederationNodesSyncConstraints(ctx, res) + if err != nil { + return err + } + + err = s.config.ErrorHandler(s.execUpsertFederationNodesSyncs(ctx, s.internalFederationNodesSyncEncoder(res))) + if err != nil { + return err + } + } + + return nil +} + +// DeleteFederationNodesSync Deletes one or more rows from federation_nodes_sync table +func (s Store) DeleteFederationNodesSync(ctx context.Context, rr ...*types.NodeSync) (err error) { + for _, res := range rr { + + err = s.execDeleteFederationNodesSyncs(ctx, squirrel.Eq{ + s.preprocessColumn("fdns.node_id", ""): store.PreprocessValue(res.NodeID, ""), + }) + if err != nil { + return s.config.ErrorHandler(err) + } + } + + return nil +} + +// DeleteFederationNodesSyncByNodeID Deletes row from the federation_nodes_sync table +func (s Store) DeleteFederationNodesSyncByNodeID(ctx context.Context, nodeID uint64) error { + return s.execDeleteFederationNodesSyncs(ctx, squirrel.Eq{ + s.preprocessColumn("fdns.node_id", ""): store.PreprocessValue(nodeID, ""), + }) +} + +// TruncateFederationNodesSyncs Deletes all rows from the federation_nodes_sync table +func (s Store) TruncateFederationNodesSyncs(ctx context.Context) error { + return s.config.ErrorHandler(s.Truncate(ctx, s.federationNodesSyncTable())) +} + +// execLookupFederationNodesSync prepares FederationNodesSync query and executes it, +// returning types.NodeSync (or error) +func (s Store) execLookupFederationNodesSync(ctx context.Context, cnd squirrel.Sqlizer) (res *types.NodeSync, err error) { + var ( + row rowScanner + ) + + row, err = s.QueryRow(ctx, s.federationNodesSyncsSelectBuilder().Where(cnd)) + if err != nil { + return + } + + res, err = s.internalFederationNodesSyncRowScanner(row) + if err != nil { + return + } + + return res, nil +} + +// execCreateFederationNodesSyncs updates all matched (by cnd) rows in federation_nodes_sync with given data +func (s Store) execCreateFederationNodesSyncs(ctx context.Context, payload store.Payload) error { + return s.config.ErrorHandler(s.Exec(ctx, s.InsertBuilder(s.federationNodesSyncTable()).SetMap(payload))) +} + +// execUpdateFederationNodesSyncs updates all matched (by cnd) rows in federation_nodes_sync with given data +func (s Store) execUpdateFederationNodesSyncs(ctx context.Context, cnd squirrel.Sqlizer, set store.Payload) error { + return s.config.ErrorHandler(s.Exec(ctx, s.UpdateBuilder(s.federationNodesSyncTable("fdns")).Where(cnd).SetMap(set))) +} + +// execUpsertFederationNodesSyncs inserts new or updates matching (by-primary-key) rows in federation_nodes_sync with given data +func (s Store) execUpsertFederationNodesSyncs(ctx context.Context, set store.Payload) error { + upsert, err := s.config.UpsertBuilder( + s.config, + s.federationNodesSyncTable(), + set, + "node_id", + ) + + if err != nil { + return err + } + + return s.config.ErrorHandler(s.Exec(ctx, upsert)) +} + +// execDeleteFederationNodesSyncs Deletes all matched (by cnd) rows in federation_nodes_sync with given data +func (s Store) execDeleteFederationNodesSyncs(ctx context.Context, cnd squirrel.Sqlizer) error { + return s.config.ErrorHandler(s.Exec(ctx, s.DeleteBuilder(s.federationNodesSyncTable("fdns")).Where(cnd))) +} + +func (s Store) internalFederationNodesSyncRowScanner(row rowScanner) (res *types.NodeSync, err error) { + res = &types.NodeSync{} + + if _, has := s.config.RowScanners["federationNodesSync"]; has { + scanner := s.config.RowScanners["federationNodesSync"].(func(_ rowScanner, _ *types.NodeSync) error) + err = scanner(row, res) + } else { + err = row.Scan( + &res.NodeID, + &res.SyncType, + &res.SyncStatus, + &res.TimeOfAction, + ) + } + + if err == sql.ErrNoRows { + return nil, store.ErrNotFound + } + + if err != nil { + return nil, fmt.Errorf("could not scan db row for FederationNodesSync: %w", err) + } else { + return res, nil + } +} + +// QueryFederationNodesSyncs returns squirrel.SelectBuilder with set table and all columns +func (s Store) federationNodesSyncsSelectBuilder() squirrel.SelectBuilder { + return s.SelectBuilder(s.federationNodesSyncTable("fdns"), s.federationNodesSyncColumns("fdns")...) +} + +// federationNodesSyncTable name of the db table +func (Store) federationNodesSyncTable(aa ...string) string { + var alias string + if len(aa) > 0 { + alias = " AS " + aa[0] + } + + return "federation_nodes_sync" + alias +} + +// FederationNodesSyncColumns returns all defined table columns +// +// With optional string arg, all columns are returned aliased +func (Store) federationNodesSyncColumns(aa ...string) []string { + var alias string + if len(aa) > 0 { + alias = aa[0] + "." + } + + return []string{ + alias + "node_id", + alias + "sync_type", + alias + "sync_status", + alias + "time_action", + } +} + +// {true true true true true} + +// sortableFederationNodesSyncColumns returns all FederationNodesSync columns flagged as sortable +// +// With optional string arg, all columns are returned aliased +func (Store) sortableFederationNodesSyncColumns() []string { + return []string{ + "node_id", + "time_action", + } +} + +// internalFederationNodesSyncEncoder encodes fields from types.NodeSync to store.Payload (map) +// +// Encoding is done by using generic approach or by calling encodeFederationNodesSync +// func when rdbms.customEncoder=true +func (s Store) internalFederationNodesSyncEncoder(res *types.NodeSync) store.Payload { + return store.Payload{ + "node_id": res.NodeID, + "sync_type": res.SyncType, + "sync_status": res.SyncStatus, + "time_action": res.TimeOfAction, + } +} + +// collectFederationNodesSyncCursorValues 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) +func (s Store) collectFederationNodesSyncCursorValues(res *types.NodeSync, cc ...string) *filter.PagingCursor { + var ( + cursor = &filter.PagingCursor{} + + hasUnique bool + + // All known primary key columns + + pkNode_id bool + + collect = func(cc ...string) { + for _, c := range cc { + switch c { + case "node_id": + cursor.Set(c, res.NodeID, false) + + pkNode_id = true + case "time_action": + cursor.Set(c, res.TimeOfAction, false) + + } + } + } + ) + + collect(cc...) + if !hasUnique || !(pkNode_id && true) { + collect("node_id") + } + + return cursor +} + +// checkFederationNodesSyncConstraints 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 can not rely +// on the full support (MySQL does not support conditional indexes) +func (s *Store) checkFederationNodesSyncConstraints(ctx context.Context, res *types.NodeSync) error { + // Consider resource valid when all fields in unique constraint check lookups + // have valid (non-empty) value + // + // Only string and uint64 are supported for now + // feel free to add additional types if needed + var valid = true + + if !valid { + return nil + } + + return nil +} diff --git a/store/rdbms/federation_nodes_sync.go b/store/rdbms/federation_nodes_sync.go new file mode 100644 index 000000000..ae28e938a --- /dev/null +++ b/store/rdbms/federation_nodes_sync.go @@ -0,0 +1,24 @@ +package rdbms + +import ( + "github.com/Masterminds/squirrel" + "github.com/cortezaproject/corteza-server/federation/types" +) + +func (s Store) convertFederationNodesSyncFilter(f types.NodeSyncFilter) (query squirrel.SelectBuilder, err error) { + query = s.federationNodesSyncsSelectBuilder() + + if f.NodeID > 0 { + query = query.Where("fdns.rel_node = ?", f.NodeID) + } + + if f.SyncStatus != "" { + query = query.Where("fdns.rel_compose_module = ?", f.SyncStatus) + } + + if f.SyncType != "" { + query = query.Where("fdns.rel_compose_namespace = ?", f.SyncType) + } + + return +} diff --git a/store/rdbms/rdbms_schema.go b/store/rdbms/rdbms_schema.go index 75d675b0a..625c0990f 100644 --- a/store/rdbms/rdbms_schema.go +++ b/store/rdbms/rdbms_schema.go @@ -78,6 +78,7 @@ func (s Schema) Tables() []*Table { s.FederationModuleExposed(), s.FederationModuleMapping(), s.FederationNodes(), + s.FederationNodesSync(), } } @@ -542,3 +543,12 @@ func (Schema) FederationNodes() *Table { CUDUsers, ) } + +func (Schema) FederationNodesSync() *Table { + return TableDef("federation_nodes_sync", + ColumnDef("node_id", ColumnTypeIdentifier), + ColumnDef("sync_type", ColumnTypeText), + ColumnDef("sync_status", ColumnTypeText), + ColumnDef("time_action", ColumnTypeTimestamp), + ) +} diff --git a/store/tests/gen_test.go b/store/tests/gen_test.go index ebc323d00..6f7d91175 100644 --- a/store/tests/gen_test.go +++ b/store/tests/gen_test.go @@ -17,6 +17,7 @@ package tests // - store/federation_exposed_modules.yaml // - store/federation_module_mappings.yaml // - store/federation_nodes.yaml +// - store/federation_nodes_sync.yaml // - store/federation_shared_modules.yaml // - store/labels.yaml // - store/messaging_attachments.yaml @@ -120,6 +121,11 @@ func testAllGenerated(t *testing.T, s store.Storer) { testFederationNodes(t, s) }) + // Run generated tests for FederationNodesSync + t.Run("FederationNodesSync", func(t *testing.T) { + testFederationNodesSync(t, s) + }) + // Run generated tests for FederationSharedModules t.Run("FederationSharedModules", func(t *testing.T) { testFederationSharedModules(t, s)