Added nodes sync service and status db table
This commit is contained in:
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
)
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
types:
|
||||
Node:
|
||||
NodeSync:
|
||||
noIdField: true
|
||||
ExposedModule:
|
||||
SharedModule:
|
||||
ModuleMapping:
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
@@ -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 }
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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
|
||||
}
|
||||
@@ -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),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user