Remove handling of duplicate storage diffs in watcher

- Can push this responsibility down to the transformers
- Update docs to reflect that transformers should handle duplicates
This commit is contained in:
Rob Mulholand
2019-05-20 13:29:09 -05:00
parent 79765c7998
commit e2909797fc
6 changed files with 5 additions and 30 deletions
+2 -2
View File
@@ -78,7 +78,7 @@ func (storageWatcher StorageWatcher) processRow(row utils.StorageDiffRow) {
return
}
executeErr := storageTransformer.Execute(row)
if executeErr != nil && executeErr != utils.ErrRowExists {
if executeErr != nil {
logrus.Warn(fmt.Sprintf("error executing storage transformer: %s", executeErr))
queueErr := storageWatcher.Queue.Add(row)
if queueErr != nil {
@@ -100,7 +100,7 @@ func (storageWatcher StorageWatcher) processQueue() {
continue
}
executeErr := storageTransformer.Execute(row)
if executeErr == nil || executeErr == utils.ErrRowExists {
if executeErr == nil {
storageWatcher.deleteRow(row.Id)
}
}
@@ -109,18 +109,6 @@ var _ = Describe("Storage Watcher", func() {
close(done)
})
It("does not queue row if transformer execution fails because row already exists", func(done Done) {
mockTransformer.ExecuteErr = utils.ErrRowExists
go storageWatcher.Execute(rows, errs, time.Hour)
Expect(<-errs).To(BeNil())
Consistently(func() bool {
return mockQueue.AddCalled
}).Should(BeFalse())
close(done)
})
It("queues row for later processing if transformer execution fails", func(done Done) {
mockTransformer.ExecuteErr = fakes.FakeError
@@ -199,17 +187,6 @@ var _ = Describe("Storage Watcher", func() {
close(done)
})
It("deletes row from queue if transformer execution errors because row already exists", func(done Done) {
mockTransformer.ExecuteErr = utils.ErrRowExists
go storageWatcher.Execute(rows, errs, time.Nanosecond)
Eventually(func() int {
return mockQueue.DeletePassedId
}).Should(Equal(row.Id))
close(done)
})
It("logs error if deleting persisted row fails", func(done Done) {
mockQueue.DeleteErr = fakes.FakeError
tempFile, fileErr := ioutil.TempFile("", "log")