Return error when no logs/headers available

- Replaces bool and moots question of error/bool ordering
- Also make event watcher execution synchronous
This commit is contained in:
Rob Mulholand
2019-09-18 20:55:15 -05:00
parent 2b798e00e0
commit 4fa19be90a
9 changed files with 137 additions and 164 deletions
+18 -14
View File
@@ -74,49 +74,53 @@ func (watcher *EventWatcher) AddTransformers(initializers []transformer.EventTra
}
// Extracts and delegates watched log events.
func (watcher *EventWatcher) Execute(recheckHeaders constants.TransformerExecution, errsChan chan error) {
extractErrsChan := make(chan error)
func (watcher *EventWatcher) Execute(recheckHeaders constants.TransformerExecution) error {
delegateErrsChan := make(chan error)
extractErrsChan := make(chan error)
defer close(delegateErrsChan)
defer close(extractErrsChan)
go watcher.extractLogs(recheckHeaders, extractErrsChan)
go watcher.delegateLogs(delegateErrsChan)
for {
select {
case extractErr := <-extractErrsChan:
logrus.Errorf("error extracting logs in event watcher: %s", extractErr.Error())
errsChan <- extractErr
case delegateErr := <-delegateErrsChan:
logrus.Errorf("error delegating logs in event watcher: %s", delegateErr.Error())
errsChan <- delegateErr
return delegateErr
case extractErr := <-extractErrsChan:
logrus.Errorf("error extracting logs in event watcher: %s", extractErr.Error())
return extractErr
}
}
}
func (watcher *EventWatcher) extractLogs(recheckHeaders constants.TransformerExecution, errs chan error) {
err, uncheckedHeadersFound := watcher.LogExtractor.ExtractLogs(recheckHeaders)
if err != nil {
err := watcher.LogExtractor.ExtractLogs(recheckHeaders)
if err != nil && err != logs.ErrNoUncheckedHeaders {
errs <- err
return
}
if uncheckedHeadersFound {
if err == logs.ErrNoUncheckedHeaders {
time.Sleep(NoNewDataPause)
watcher.extractLogs(recheckHeaders, errs)
} else {
time.Sleep(NoNewDataPause)
watcher.extractLogs(recheckHeaders, errs)
}
}
func (watcher *EventWatcher) delegateLogs(errs chan error) {
err, logsFound := watcher.LogDelegator.DelegateLogs()
if err != nil {
err := watcher.LogDelegator.DelegateLogs()
if err != nil && err != logs.ErrNoLogs {
errs <- err
return
}
if logsFound {
if err == logs.ErrNoLogs {
time.Sleep(NoNewDataPause)
watcher.delegateLogs(errs)
} else {
time.Sleep(NoNewDataPause)
watcher.delegateLogs(errs)
}
}
+38 -61
View File
@@ -17,6 +17,7 @@
package watcher_test
import (
"errors"
. "github.com/onsi/ginkgo"
. "github.com/onsi/gomega"
"github.com/vulcanize/vulcanizedb/libraries/shared/constants"
@@ -26,6 +27,8 @@ import (
"github.com/vulcanize/vulcanizedb/pkg/fakes"
)
var errExecuteClosed = errors.New("this error means the mocks were finished executing")
var _ = Describe("Event Watcher", func() {
var (
delegator *mocks.MockLogDelegator
@@ -57,7 +60,8 @@ var _ = Describe("Event Watcher", func() {
fakeTransformerTwo.FakeTransformerInitializer,
}
eventWatcher.AddTransformers(initializers)
err := eventWatcher.AddTransformers(initializers)
Expect(err).NotTo(HaveOccurred())
})
It("adds initialized transformer to log delegator", func() {
@@ -78,114 +82,87 @@ var _ = Describe("Event Watcher", func() {
})
Describe("Execute", func() {
var errsChan chan error
BeforeEach(func() {
errsChan = make(chan error)
})
It("extracts watched logs", func(done Done) {
delegator.DelegateErrors = []error{nil}
delegator.LogsFound = []bool{false}
extractor.ExtractLogsErrors = []error{nil}
extractor.UncheckedHeadersExist = []bool{false}
extractor.ExtractLogsErrors = []error{nil, errExecuteClosed}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Eventually(func() int {
return extractor.ExtractLogsCount
}).Should(Equal(1))
Expect(err).To(MatchError(errExecuteClosed))
Eventually(func() bool {
return extractor.ExtractLogsCount > 0
}).Should(BeTrue())
close(done)
})
It("returns error if extracting logs fails", func(done Done) {
delegator.DelegateErrors = []error{nil}
delegator.LogsFound = []bool{false}
extractor.ExtractLogsErrors = []error{fakes.FakeError}
extractor.UncheckedHeadersExist = []bool{false}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Expect(<-errsChan).To(MatchError(fakes.FakeError))
Expect(err).To(MatchError(fakes.FakeError))
close(done)
})
It("extracts watched logs again if missing headers found", func(done Done) {
delegator.DelegateErrors = []error{nil}
delegator.LogsFound = []bool{false}
extractor.ExtractLogsErrors = []error{nil, nil}
extractor.UncheckedHeadersExist = []bool{true, false}
extractor.ExtractLogsErrors = []error{nil, errExecuteClosed}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Eventually(func() int {
return extractor.ExtractLogsCount
}).Should(Equal(2))
Expect(err).To(MatchError(errExecuteClosed))
Eventually(func() bool {
return extractor.ExtractLogsCount > 1
}).Should(BeTrue())
close(done)
})
It("returns error if extracting logs fails on subsequent run", func(done Done) {
delegator.DelegateErrors = []error{nil}
delegator.LogsFound = []bool{false}
extractor.ExtractLogsErrors = []error{nil, fakes.FakeError}
extractor.UncheckedHeadersExist = []bool{true, false}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Expect(<-errsChan).To(MatchError(fakes.FakeError))
Expect(err).To(MatchError(fakes.FakeError))
close(done)
})
It("delegates untransformed logs", func(done Done) {
delegator.DelegateErrors = []error{nil}
delegator.LogsFound = []bool{false}
extractor.ExtractLogsErrors = []error{nil}
extractor.UncheckedHeadersExist = []bool{false}
It("delegates untransformed logs", func() {
delegator.DelegateErrors = []error{nil, errExecuteClosed}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Eventually(func() int {
return delegator.DelegateCallCount
}).Should(Equal(1))
close(done)
Expect(err).To(MatchError(errExecuteClosed))
Eventually(func() bool {
return delegator.DelegateCallCount > 0
}).Should(BeTrue())
})
It("returns error if delegating logs fails", func(done Done) {
delegator.LogsFound = []bool{false}
delegator.DelegateErrors = []error{fakes.FakeError}
extractor.ExtractLogsErrors = []error{nil}
extractor.UncheckedHeadersExist = []bool{false}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Expect(<-errsChan).To(MatchError(fakes.FakeError))
Expect(err).To(MatchError(fakes.FakeError))
close(done)
})
It("delegates logs again if untransformed logs found", func(done Done) {
delegator.DelegateErrors = []error{nil, nil}
delegator.LogsFound = []bool{true, false}
extractor.ExtractLogsErrors = []error{nil}
extractor.UncheckedHeadersExist = []bool{false}
delegator.DelegateErrors = []error{nil, nil, nil, errExecuteClosed}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Eventually(func() int {
return delegator.DelegateCallCount
}).Should(Equal(2))
Expect(err).To(MatchError(errExecuteClosed))
Eventually(func() bool {
return delegator.DelegateCallCount > 1
}).Should(BeTrue())
close(done)
})
It("returns error if delegating logs fails on subsequent run", func(done Done) {
delegator.DelegateErrors = []error{nil, fakes.FakeError}
delegator.LogsFound = []bool{true, false}
extractor.ExtractLogsErrors = []error{nil}
extractor.UncheckedHeadersExist = []bool{false}
go eventWatcher.Execute(constants.HeaderUnchecked, errsChan)
err := eventWatcher.Execute(constants.HeaderUnchecked)
Expect(<-errsChan).To(MatchError(fakes.FakeError))
Expect(err).To(MatchError(fakes.FakeError))
close(done)
})
})