forked from cerc-io/ipld-eth-server
Use new config getter on shared.Transformer <3
This commit is contained in:
+16
-16
@@ -24,7 +24,8 @@ type Watcher struct {
|
||||
Repository WatcherRepository
|
||||
}
|
||||
|
||||
func NewWatcher(db *postgres.DB, fetcher shared.LogFetcher, repository WatcherRepository, chunker shared.Chunker) Watcher {
|
||||
func NewWatcher(db *postgres.DB, fetcher shared.LogFetcher, repository WatcherRepository) Watcher {
|
||||
chunker := shared.NewLogChunker()
|
||||
return Watcher{
|
||||
DB: db,
|
||||
Fetcher: fetcher,
|
||||
@@ -33,25 +34,21 @@ func NewWatcher(db *postgres.DB, fetcher shared.LogFetcher, repository WatcherRe
|
||||
}
|
||||
}
|
||||
|
||||
// Adds transformers to the watcher, each needs an initializer and the associated config.
|
||||
// This also changes the configuration of the chunker, so that it will consider the new transformers.
|
||||
func (watcher *Watcher) AddTransformers(initializers []shared.TransformerInitializer, configs []shared.TransformerConfig) {
|
||||
if len(initializers) != len(configs) {
|
||||
panic("Mismatch in number of transformers initializers and configs!")
|
||||
}
|
||||
// Adds transformers to the watcher and updates the chunker, so that it will consider the new transformers.
|
||||
func (watcher *Watcher) AddTransformers(initializers []shared.TransformerInitializer) {
|
||||
var contractAddresses []common.Address
|
||||
var topic0s []common.Hash
|
||||
var configs []shared.TransformerConfig
|
||||
|
||||
for _, initializer := range initializers {
|
||||
transformer := initializer(watcher.DB)
|
||||
watcher.Transformers = append(watcher.Transformers, transformer)
|
||||
}
|
||||
|
||||
var contractAddresses []common.Address
|
||||
var topic0s []common.Hash
|
||||
config := transformer.GetConfig()
|
||||
configs = append(configs, config)
|
||||
|
||||
for _, config := range configs {
|
||||
for _, address := range config.ContractAddresses {
|
||||
contractAddresses = append(contractAddresses, common.HexToAddress(address))
|
||||
}
|
||||
addresses := shared.HexStringsToAddresses(config.ContractAddresses)
|
||||
contractAddresses = append(contractAddresses, addresses...)
|
||||
topic0s = append(topic0s, common.HexToHash(config.Topic))
|
||||
}
|
||||
|
||||
@@ -85,11 +82,14 @@ func (watcher *Watcher) Execute() error {
|
||||
|
||||
chunkedLogs := watcher.Chunker.ChunkLogs(logs)
|
||||
|
||||
// Can't quit early and mark as checked if there are no logs. If we are running continuousLogSync,
|
||||
// not all logs we're interested in might have been fetched.
|
||||
for _, transformer := range watcher.Transformers {
|
||||
logChunk := chunkedLogs[transformer.Name()]
|
||||
transformerName := transformer.GetConfig().TransformerName
|
||||
logChunk := chunkedLogs[transformerName]
|
||||
err = transformer.Execute(logChunk, header)
|
||||
if err != nil {
|
||||
log.Error("%v transformer failed to execute in watcher: %v", transformer.Name(), err)
|
||||
log.Error("%v transformer failed to execute in watcher: %v", transformerName, err)
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -21,7 +21,7 @@ type MockTransformer struct {
|
||||
executeError error
|
||||
passedLogs []types.Log
|
||||
passedHeader core.Header
|
||||
transformerName string
|
||||
config shared2.TransformerConfig
|
||||
}
|
||||
|
||||
func (mh *MockTransformer) Execute(logs []types.Log, header core.Header) error {
|
||||
@@ -34,57 +34,60 @@ func (mh *MockTransformer) Execute(logs []types.Log, header core.Header) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (mh *MockTransformer) Name() string {
|
||||
return mh.transformerName
|
||||
func (mh *MockTransformer) GetConfig() shared2.TransformerConfig {
|
||||
return mh.config
|
||||
}
|
||||
|
||||
func (mh *MockTransformer) SetTransformerName(name string) {
|
||||
mh.transformerName = name
|
||||
func (mh *MockTransformer) SetTransformerConfig(config shared2.TransformerConfig) {
|
||||
mh.config = config
|
||||
}
|
||||
|
||||
func fakeTransformerInitializer(db *postgres.DB) shared2.Transformer {
|
||||
return &MockTransformer{}
|
||||
func (mh *MockTransformer) FakeTransformerInitializer(db *postgres.DB) shared2.Transformer {
|
||||
return mh
|
||||
}
|
||||
|
||||
var fakeTransformerConfig = []shared2.TransformerConfig{{
|
||||
var fakeTransformerConfig = shared2.TransformerConfig{
|
||||
TransformerName: "FakeTransformer",
|
||||
ContractAddresses: []string{"FakeAddress"},
|
||||
Topic: "FakeTopic",
|
||||
}}
|
||||
}
|
||||
|
||||
var _ = Describe("Watcher", func() {
|
||||
It("initialises correctly", func() {
|
||||
db := test_config.NewTestDB(core.Node{ID: "testNode"})
|
||||
fetcher := &mocks.MockLogFetcher{}
|
||||
repository := &mocks.MockWatcherRepository{}
|
||||
chunker := shared2.NewLogChunker()
|
||||
|
||||
watcher := shared.NewWatcher(db, fetcher, repository, chunker)
|
||||
watcher := shared.NewWatcher(db, fetcher, repository)
|
||||
|
||||
Expect(watcher.DB).To(Equal(db))
|
||||
Expect(watcher.Fetcher).To(Equal(fetcher))
|
||||
Expect(watcher.Chunker).To(Equal(chunker))
|
||||
Expect(watcher.Chunker).NotTo(BeNil())
|
||||
Expect(watcher.Repository).To(Equal(repository))
|
||||
})
|
||||
|
||||
It("adds transformers", func() {
|
||||
chunker := shared2.NewLogChunker()
|
||||
watcher := shared.NewWatcher(nil, nil, nil, chunker)
|
||||
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformerInitializer}, fakeTransformerConfig)
|
||||
watcher := shared.NewWatcher(nil, nil, nil)
|
||||
fakeTransformer := &MockTransformer{}
|
||||
fakeTransformer.SetTransformerConfig(fakeTransformerConfig)
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformer.FakeTransformerInitializer})
|
||||
|
||||
Expect(len(watcher.Transformers)).To(Equal(1))
|
||||
Expect(watcher.Transformers).To(ConsistOf(&MockTransformer{}))
|
||||
Expect(watcher.Transformers).To(ConsistOf(fakeTransformer))
|
||||
Expect(watcher.Topics).To(Equal([]common.Hash{common.HexToHash("FakeTopic")}))
|
||||
Expect(watcher.Addresses).To(Equal([]common.Address{common.HexToAddress("FakeAddress")}))
|
||||
})
|
||||
|
||||
It("adds transformers from multiple sources", func() {
|
||||
chunker := shared2.NewLogChunker()
|
||||
watcher := shared.NewWatcher(nil, nil, nil, chunker)
|
||||
watcher := shared.NewWatcher(nil, nil, nil)
|
||||
fakeTransformer1 := &MockTransformer{}
|
||||
fakeTransformer1.SetTransformerConfig(fakeTransformerConfig)
|
||||
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformerInitializer}, fakeTransformerConfig)
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformerInitializer}, fakeTransformerConfig)
|
||||
fakeTransformer2 := &MockTransformer{}
|
||||
fakeTransformer2.SetTransformerConfig(fakeTransformerConfig)
|
||||
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformer1.FakeTransformerInitializer})
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{fakeTransformer2.FakeTransformerInitializer})
|
||||
|
||||
Expect(len(watcher.Transformers)).To(Equal(2))
|
||||
Expect(watcher.Topics).To(Equal([]common.Hash{common.HexToHash("FakeTopic"),
|
||||
@@ -101,7 +104,6 @@ var _ = Describe("Watcher", func() {
|
||||
headerRepository repositories.HeaderRepository
|
||||
mockFetcher mocks.MockLogFetcher
|
||||
repository mocks.MockWatcherRepository
|
||||
chunker *shared2.LogChunker
|
||||
)
|
||||
|
||||
BeforeEach(func() {
|
||||
@@ -111,10 +113,9 @@ var _ = Describe("Watcher", func() {
|
||||
headerRepository = repositories.NewHeaderRepository(db)
|
||||
_, err := headerRepository.CreateOrUpdateHeader(fakes.FakeHeader)
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
chunker = shared2.NewLogChunker()
|
||||
|
||||
repository = mocks.MockWatcherRepository{}
|
||||
watcher = shared.NewWatcher(db, &mockFetcher, &repository, chunker)
|
||||
watcher = shared.NewWatcher(db, &mockFetcher, &repository)
|
||||
})
|
||||
|
||||
It("executes each transformer", func() {
|
||||
@@ -141,9 +142,7 @@ var _ = Describe("Watcher", func() {
|
||||
|
||||
It("passes only relevant logs to each transformer", func() {
|
||||
transformerA := &MockTransformer{}
|
||||
transformerA.SetTransformerName("transformerA")
|
||||
transformerB := &MockTransformer{}
|
||||
transformerB.SetTransformerName("transformerB")
|
||||
|
||||
configA := shared2.TransformerConfig{TransformerName: "transformerA",
|
||||
ContractAddresses: []string{"0x000000000000000000000000000000000000000A"},
|
||||
@@ -151,7 +150,9 @@ var _ = Describe("Watcher", func() {
|
||||
configB := shared2.TransformerConfig{TransformerName: "transformerB",
|
||||
ContractAddresses: []string{"0x000000000000000000000000000000000000000b"},
|
||||
Topic: "0xB"}
|
||||
configs := []shared2.TransformerConfig{configA, configB}
|
||||
|
||||
transformerA.SetTransformerConfig(configA)
|
||||
transformerB.SetTransformerConfig(configB)
|
||||
|
||||
logA := types.Log{Address: common.HexToAddress("0xA"),
|
||||
Topics: []common.Hash{common.HexToHash("0xA")}}
|
||||
@@ -159,11 +160,10 @@ var _ = Describe("Watcher", func() {
|
||||
Topics: []common.Hash{common.HexToHash("0xB")}}
|
||||
mockFetcher.SetFetchedLogs([]types.Log{logA, logB})
|
||||
|
||||
chunker.AddConfigs(configs)
|
||||
|
||||
repository.SetMissingHeaders([]core.Header{fakes.FakeHeader})
|
||||
watcher = shared.NewWatcher(db, &mockFetcher, &repository, chunker)
|
||||
watcher.Transformers = []shared2.Transformer{transformerA, transformerB}
|
||||
watcher = shared.NewWatcher(db, &mockFetcher, &repository)
|
||||
watcher.AddTransformers([]shared2.TransformerInitializer{
|
||||
transformerA.FakeTransformerInitializer, transformerB.FakeTransformerInitializer})
|
||||
|
||||
err := watcher.Execute()
|
||||
Expect(err).NotTo(HaveOccurred())
|
||||
|
||||
Reference in New Issue
Block a user