diff --git a/federation/service/processer_data.go b/federation/service/processer_data.go new file mode 100644 index 000000000..1b7f8f01e --- /dev/null +++ b/federation/service/processer_data.go @@ -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 +} diff --git a/federation/service/processer_structure.go b/federation/service/processer_structure.go new file mode 100644 index 000000000..04f1a8246 --- /dev/null +++ b/federation/service/processer_structure.go @@ -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 +}