Added data, structure processers
This commit is contained in:
@@ -0,0 +1,76 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
ct "github.com/cortezaproject/corteza-server/compose/types"
|
||||
"github.com/cortezaproject/corteza-server/federation/types"
|
||||
"github.com/cortezaproject/corteza-server/pkg/auth"
|
||||
"github.com/cortezaproject/corteza-server/pkg/decoder"
|
||||
st "github.com/cortezaproject/corteza-server/system/types"
|
||||
)
|
||||
|
||||
type (
|
||||
dataProcesser struct {
|
||||
ID uint64
|
||||
ComposeModuleID uint64
|
||||
ComposeNamespaceID uint64
|
||||
NodeBaseURL string
|
||||
ModuleMappings *types.ModuleFieldMappingSet
|
||||
ModuleMappingValues *ct.RecordValueSet
|
||||
SyncService *Sync
|
||||
Node *types.Node
|
||||
User *st.User
|
||||
}
|
||||
|
||||
dataProcesserResponse struct {
|
||||
Processed int
|
||||
}
|
||||
)
|
||||
|
||||
// Process gets the payload from syncer and
|
||||
// uses the decode package to decode the whole set, depending on
|
||||
// the filtering that was used (limit)
|
||||
func (dp *dataProcesser) Process(ctx context.Context, payload []byte) (ProcesserResponse, error) {
|
||||
processed := 0
|
||||
o, err := decoder.DecodeFederationRecordSync([]byte(payload))
|
||||
|
||||
if err != nil {
|
||||
return dataProcesserResponse{
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
if len(o) == 0 {
|
||||
return dataProcesserResponse{
|
||||
Processed: processed,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// get the user that is tied to this node
|
||||
ctx = auth.SetIdentityToContext(ctx, dp.User)
|
||||
|
||||
for _, er := range o {
|
||||
dp.SyncService.mapper.Merge(&er.Values, dp.ModuleMappingValues, dp.ModuleMappings)
|
||||
|
||||
rec := &ct.Record{
|
||||
ModuleID: dp.ComposeModuleID,
|
||||
NamespaceID: dp.ComposeNamespaceID,
|
||||
Values: *dp.ModuleMappingValues,
|
||||
}
|
||||
|
||||
AddFederationLabel(rec, dp.NodeBaseURL)
|
||||
|
||||
_, err := dp.SyncService.CreateRecord(ctx, rec)
|
||||
|
||||
if err != nil {
|
||||
continue
|
||||
}
|
||||
|
||||
processed++
|
||||
}
|
||||
|
||||
return dataProcesserResponse{
|
||||
Processed: processed,
|
||||
}, nil
|
||||
}
|
||||
@@ -0,0 +1,124 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/cortezaproject/corteza-server/federation/types"
|
||||
"github.com/cortezaproject/corteza-server/pkg/auth"
|
||||
"github.com/cortezaproject/corteza-server/pkg/decoder"
|
||||
st "github.com/cortezaproject/corteza-server/system/types"
|
||||
)
|
||||
|
||||
type (
|
||||
structureProcesser struct {
|
||||
SyncService *Sync
|
||||
Node *types.Node
|
||||
User *st.User
|
||||
}
|
||||
|
||||
structureProcesserResponse struct {
|
||||
ModuleID uint64
|
||||
Processed int
|
||||
}
|
||||
)
|
||||
|
||||
// Process gets the payload from syncer and
|
||||
// uses the decode package to decode the whole set, depending on
|
||||
// the filtering that was used (limit)
|
||||
func (dp *structureProcesser) Process(ctx context.Context, payload []byte) (ProcesserResponse, error) {
|
||||
processed := 0
|
||||
o, err := decoder.DecodeFederationModuleSync([]byte(payload))
|
||||
|
||||
if err != nil {
|
||||
return structureProcesserResponse{
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
if len(o) == 0 {
|
||||
return structureProcesserResponse{
|
||||
Processed: processed,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// get the user that is tied to this node
|
||||
ctx = auth.SetIdentityToContext(ctx, dp.User)
|
||||
|
||||
for _, em := range o {
|
||||
new := &types.SharedModule{
|
||||
NodeID: dp.Node.ID,
|
||||
ExternalFederationModuleID: em.ID,
|
||||
Fields: em.Fields,
|
||||
Handle: em.Handle,
|
||||
Name: em.Name,
|
||||
}
|
||||
|
||||
existing, err := dp.SyncService.LookupSharedModule(ctx, new)
|
||||
|
||||
if err != nil {
|
||||
return structureProcesserResponse{
|
||||
ModuleID: existing.ExternalFederationModuleID,
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
// create
|
||||
{
|
||||
if existing == nil {
|
||||
_, err := dp.SyncService.CreateSharedModule(ctx, new)
|
||||
|
||||
if err != nil {
|
||||
return structureProcesserResponse{
|
||||
ModuleID: new.ExternalFederationModuleID,
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
processed++
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// update
|
||||
{
|
||||
canUpdate, err := dp.SyncService.CanUpdateSharedModule(ctx, new, existing)
|
||||
|
||||
if err != nil {
|
||||
return structureProcesserResponse{
|
||||
ModuleID: existing.ExternalFederationModuleID,
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
// stop the sync for this module
|
||||
if !canUpdate {
|
||||
return structureProcesserResponse{
|
||||
ModuleID: new.ExternalFederationModuleID,
|
||||
Processed: processed,
|
||||
}, SharedModuleErrFederationSyncStructureChanged(&sharedModuleActionProps{module: existing})
|
||||
}
|
||||
|
||||
existing.Name = new.Name
|
||||
existing.Handle = new.Handle
|
||||
existing.Fields = new.Fields
|
||||
|
||||
_, err = dp.SyncService.UpdateSharedModule(ctx, existing)
|
||||
|
||||
if err != nil {
|
||||
return structureProcesserResponse{
|
||||
ModuleID: existing.ExternalFederationModuleID,
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
|
||||
processed++
|
||||
continue
|
||||
}
|
||||
}
|
||||
|
||||
// use paging, but only one per-time
|
||||
return structureProcesserResponse{
|
||||
ModuleID: o[0].ID,
|
||||
Processed: processed,
|
||||
}, err
|
||||
}
|
||||
Reference in New Issue
Block a user