diff --git a/federation/rest/sync_data.go b/federation/rest/sync_data.go index b0b6d0b3e..68f93133e 100644 --- a/federation/rest/sync_data.go +++ b/federation/rest/sync_data.go @@ -134,6 +134,7 @@ func (ctrl SyncData) ReadExposed(ctx context.Context, r *request.SyncDataReadExp ModuleID: em.ComposeModuleID, Query: buildLastSyncQuery(r.LastSync), Check: ignoreFederated(users), + Deleted: filter.StateInclusive, } if f.Paging, err = filter.NewPaging(r.Limit, r.PageCursor); err != nil { @@ -175,7 +176,8 @@ func buildLastSyncQuery(ts uint64) string { } return fmt.Sprintf( - "(updated_at >= '%s' OR created_at >= '%s')", + "(updated_at >= '%s' OR created_at >= '%s' OR deleted_at >= '%s')", + t.UTC().Format(time.RFC3339), t.UTC().Format(time.RFC3339), t.UTC().Format(time.RFC3339)) } diff --git a/federation/service/processer_data.go b/federation/service/processer_data.go index c6302100a..7e0245580 100644 --- a/federation/service/processer_data.go +++ b/federation/service/processer_data.go @@ -58,11 +58,20 @@ func (dp *dataProcesser) Process(ctx context.Context, payload []byte) (Processer for _, er := range o { dp.SyncService.mapper.Merge(&er.Values, dp.ModuleMappingValues, dp.ModuleMappings) + if er.DeletedAt != nil { + // find the record + if rec, err = dp.findRecordByFederationID(ctx, er.ID, dp.ComposeModuleID, dp.ComposeNamespaceID); err != nil { + continue + } + + dp.SyncService.DeleteRecord(ctx, rec) + processed++ + + continue + } + if er.UpdatedAt != nil { - - rec, err = dp.findRecordByFederationID(ctx, er.ID, dp.ComposeModuleID, dp.ComposeNamespaceID) - - if err != nil { + if rec, err = dp.findRecordByFederationID(ctx, er.ID, dp.ComposeModuleID, dp.ComposeNamespaceID); err != nil { // could not find existing record continue } @@ -75,7 +84,6 @@ func (dp *dataProcesser) Process(ctx context.Context, payload []byte) (Processer // if the record was updated on origin, but we somehow do not have it // create it anyway if rec == nil { - rec = &ct.Record{ ModuleID: dp.ComposeModuleID, NamespaceID: dp.ComposeNamespaceID, @@ -84,7 +92,6 @@ func (dp *dataProcesser) Process(ctx context.Context, payload []byte) (Processer AddFederationLabel(rec, "federation", dp.NodeBaseURL) AddFederationLabel(rec, "federation_extrecord", fmt.Sprintf("%d", er.ID)) - } if rec.ID != 0 { diff --git a/federation/service/processer_data_test.go b/federation/service/processer_data_test.go index d712c1cc3..0e996a7dd 100644 --- a/federation/service/processer_data_test.go +++ b/federation/service/processer_data_test.go @@ -28,6 +28,9 @@ type ( testRecordServiceUpdateSuccess struct { cs.RecordService } + testRecordServiceDeleteSuccess struct { + cs.RecordService + } testRecordServicePersistError struct { cs.RecordService } @@ -81,6 +84,21 @@ func TestProcesserData_persist(t *testing.T) { &testUserService{}, &testRoleService{}), }, + { + "successful delete on valid mapping and existing federated record", + `{"response": {"set": [{"recordID":"1","values":[{"name":"Facebook","value":"foobar"}],"createdAt":"2020-12-05T10:10:10Z", "deletedAt":"2020-12-07T10:10:10Z"}]}}`, + `[{"origin":{"kind":"String","name":"Description","label":"Description","isMulti":false},"destination":{"kind":"String","name":"Name","label":"Description","isMulti":false}},{"origin":{"kind":"Url","name":"Facebook","label":"Facebook","isMulti":false},"destination":{"kind":"Url","name":"Fb","label":"Facebook","isMulti":false}}]`, + 1, + "", + &ct.RecordValueSet{&ct.RecordValue{Name: "Fb", Value: ""}}, + NewSync( + &Syncer{}, + &Mapper{}, + &testSharedModuleService{}, + &testRecordServiceDeleteSuccess{}, + &testUserService{}, + &testRoleService{}), + }, { "persist error on valid mapping", `{"response": {"set": [{"recordID":"1","values":[{"name":"Facebook","value":"foobar"}]}]}}`, @@ -205,6 +223,19 @@ func (s testRecordServiceUpdateSuccess) With(_ context.Context) cs.RecordService return &testRecordServiceUpdateSuccess{} } +// delete success +func (s testRecordServiceDeleteSuccess) DeleteByID(namespaceID, moduleID uint64, recordID ...uint64) error { + return nil +} + +func (s testRecordServiceDeleteSuccess) Find(filter ct.RecordFilter) (ct.RecordSet, ct.RecordFilter, error) { + return ct.RecordSet{&ct.Record{ID: 2, ModuleID: 2, NamespaceID: 2}}, ct.RecordFilter{}, nil +} + +func (s testRecordServiceDeleteSuccess) With(_ context.Context) cs.RecordService { + return &testRecordServiceDeleteSuccess{} +} + // create error func (s testRecordServicePersistError) Create(record *ct.Record) (*ct.Record, error) { return nil, errors.New("mocked error") diff --git a/federation/service/sync.go b/federation/service/sync.go index b1a9777ad..081af1e93 100644 --- a/federation/service/sync.go +++ b/federation/service/sync.go @@ -85,6 +85,11 @@ func (s *Sync) UpdateRecord(ctx context.Context, rec *ct.Record) (*ct.Record, er return s.composeRecordService.With(ctx).Update(rec) } +// DeleteRecord wraps the compose Record service Update +func (s *Sync) DeleteRecord(ctx context.Context, rec *ct.Record) error { + return s.composeRecordService.With(ctx).DeleteByID(rec.NamespaceID, rec.ModuleID, rec.ID) +} + // FindRecord find the record via federation label func (s *Sync) FindRecords(ctx context.Context, filter ct.RecordFilter) (set ct.RecordSet, err error) { set, _, err = s.composeRecordService.With(ctx).Find(filter)