Added delete functionality
This commit is contained in:
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user