Remove unnecessary async from the event watcher

- extract and delegate logs synchronously after initial goroutine fired
This commit is contained in:
Rob Mulholand
2019-08-28 09:25:13 -05:00
parent 1883a11ab1
commit d76be4962b
7 changed files with 149 additions and 217 deletions
+15 -16
View File
@@ -30,9 +30,14 @@ import (
var ErrNoWatchedAddresses = errors.New("no watched addresses configured in the log extractor")
const (
missingHeadersFound = true
noMissingHeadersFound = false
)
type ILogExtractor interface {
AddTransformerConfig(config transformer.EventTransformerConfig)
ExtractLogs(recheckHeaders constants.TransformerExecution, errs chan error, missingHeadersFound chan bool)
ExtractLogs(recheckHeaders constants.TransformerExecution) (error, bool)
}
type LogExtractor struct {
@@ -59,56 +64,50 @@ func (extractor *LogExtractor) AddTransformerConfig(config transformer.EventTran
}
// Fetch and persist watched logs
func (extractor LogExtractor) ExtractLogs(recheckHeaders constants.TransformerExecution, errs chan error, missingHeadersFound chan bool) {
func (extractor LogExtractor) ExtractLogs(recheckHeaders constants.TransformerExecution) (error, bool) {
if len(extractor.Addresses) < 1 {
logrus.Errorf("error extracting logs: %s", ErrNoWatchedAddresses.Error())
errs <- ErrNoWatchedAddresses
return
return ErrNoWatchedAddresses, noMissingHeadersFound
}
missingHeaders, missingHeadersErr := extractor.CheckedHeadersRepository.MissingHeaders(*extractor.StartingBlock, -1, getCheckCount(recheckHeaders))
if missingHeadersErr != nil {
logrus.Errorf("error fetching missing headers: %s", missingHeadersErr)
errs <- missingHeadersErr
return
return missingHeadersErr, noMissingHeadersFound
}
if len(missingHeaders) < 1 {
missingHeadersFound <- false
return
return nil, noMissingHeadersFound
}
for _, header := range missingHeaders {
logs, fetchLogsErr := extractor.Fetcher.FetchLogs(extractor.Addresses, extractor.Topics, header)
if fetchLogsErr != nil {
logError("error fetching logs for header: %s", fetchLogsErr, header)
errs <- fetchLogsErr
return
return fetchLogsErr, missingHeadersFound
}
if len(logs) > 0 {
transactionsSyncErr := extractor.Syncer.SyncTransactions(header.Id, logs)
if transactionsSyncErr != nil {
logError("error syncing transactions: %s", transactionsSyncErr, header)
errs <- transactionsSyncErr
return
return transactionsSyncErr, missingHeadersFound
}
createLogsErr := extractor.LogRepository.CreateHeaderSyncLogs(header.Id, logs)
if createLogsErr != nil {
logError("error persisting logs: %s", createLogsErr, header)
errs <- createLogsErr
return
return createLogsErr, missingHeadersFound
}
}
markHeaderCheckedErr := extractor.CheckedHeadersRepository.MarkHeaderChecked(header.Id)
if markHeaderCheckedErr != nil {
logError("error marking header checked: %s", markHeaderCheckedErr, header)
errs <- markHeaderCheckedErr
return markHeaderCheckedErr, missingHeadersFound
}
}
missingHeadersFound <- true
return nil, missingHeadersFound
}
func earlierStartingBlockNumber(transformerBlock, watcherBlock int64) bool {