diff --git a/config.md b/config.md index 4780250..cac0c78 100644 --- a/config.md +++ b/config.md @@ -344,6 +344,7 @@ |---|-----------|----|-------------| |address|Listener address|`int`|`127.0.0.1` |enabled|Enables the monitoring APIs|`boolean`|`false` +|loggingPath|The path at which to serve the dynamic logging API, which allows the log level to be changed at runtime|`string`|`/logging` |metricsPath|The path from which to serve the Prometheus metrics|`string`|`/metrics` |port|Listener port|`int`|`6000` |publicURL|Externally available URL for the HTTP endpoint|`string`|`` diff --git a/go.mod b/go.mod index d00d1d0..72388ac 100644 --- a/go.mod +++ b/go.mod @@ -9,8 +9,8 @@ require ( github.com/hashicorp/golang-lru v1.0.2 github.com/hyperledger-firefly/common v1.6.5 github.com/hyperledger-firefly/signer v1.2.1 - github.com/hyperledger-firefly/transaction-manager v1.5.2 - github.com/sirupsen/logrus v1.10.0 + github.com/hyperledger-firefly/transaction-manager v1.5.3 + github.com/sirupsen/logrus v1.10.1 github.com/spf13/cobra v1.10.2 github.com/stretchr/testify v1.12.1 golang.org/x/net v0.58.0 @@ -32,7 +32,7 @@ require ( github.com/decred/dcrd/dcrec/secp256k1/v4 v4.2.0 // indirect github.com/docker/go-units v0.5.0 // indirect github.com/fsnotify/fsnotify v1.9.0 // indirect - github.com/getkin/kin-openapi v0.144.0 // indirect + github.com/getkin/kin-openapi v0.147.0 // indirect github.com/ghodss/yaml v1.0.0 // indirect github.com/go-openapi/jsonpointer v0.22.5 // indirect github.com/go-openapi/swag/jsonname v0.25.5 // indirect @@ -57,12 +57,12 @@ require ( github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 // indirect github.com/oasdiff/yaml v0.1.1 // indirect github.com/oasdiff/yaml3 v0.0.14 // indirect - github.com/oklog/ulid/v2 v2.1.1 // indirect + github.com/oklog/ulid/v2 v2.1.2 // indirect github.com/pelletier/go-toml/v2 v2.2.4 // indirect github.com/pkg/errors v0.9.1 // indirect - github.com/prometheus/client_golang v1.24.0 // indirect + github.com/prometheus/client_golang v1.24.1 // indirect github.com/prometheus/client_model v0.6.2 // indirect - github.com/prometheus/common v0.70.0 // indirect + github.com/prometheus/common v0.70.1 // indirect github.com/prometheus/procfs v0.21.1 // indirect github.com/rs/cors v1.11.1 // indirect github.com/sagikazarmark/locafero v0.11.0 // indirect diff --git a/go.sum b/go.sum index 998dac9..53879e2 100644 --- a/go.sum +++ b/go.sum @@ -56,8 +56,8 @@ github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMo github.com/fsnotify/fsnotify v1.4.9/go.mod h1:znqG4EE+3YCdAaPaxE2ZRY/06pZUdp0tY4IgpuI1SZQ= github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= -github.com/getkin/kin-openapi v0.144.0 h1:hIRcTH+KjLfkLpYU6bSSfdFpi0fZi1fp+hSPi4aQu9Y= -github.com/getkin/kin-openapi v0.144.0/go.mod h1:3BH9M9XDe/y9M5DSvEocVYAYq1w0qrhJHjC/vZi0AaY= +github.com/getkin/kin-openapi v0.147.0 h1:s+Xsm9gUMPJbgCnABZ2to3zSQQ5A9dyj/zo62VVsldY= +github.com/getkin/kin-openapi v0.147.0/go.mod h1:3BH9M9XDe/y9M5DSvEocVYAYq1w0qrhJHjC/vZi0AaY= github.com/ghodss/yaml v1.0.0 h1:wQHKEahhL6wmXdzwWG11gIVCkOv05bNOh+Rxn0yngAk= github.com/ghodss/yaml v1.0.0/go.mod h1:4dBDuWmgqj2HViK6kFavaiC9ZROes6MMH2rRYeMEF04= github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= @@ -109,8 +109,8 @@ github.com/hyperledger-firefly/common v1.6.5 h1:/LEY1LlXQEezLVnEzSj8Lq3WQ9Uh1qUX github.com/hyperledger-firefly/common v1.6.5/go.mod h1:h1LhJHfJYH7hTioyx0BsjSt5sDNsxhXVV1jDaKaBszU= github.com/hyperledger-firefly/signer v1.2.1 h1:YmnCfiOPhrP/6WIGt5KwV3wiDEI61k7h282MpagSNnI= github.com/hyperledger-firefly/signer v1.2.1/go.mod h1:qLBLnXZZHk83x68WLvDjx8q80+6lQzx57FPeRyI6jmo= -github.com/hyperledger-firefly/transaction-manager v1.5.2 h1:tOqwKD6VDHKTkwdCQTTSP1Vuhxg8DOwsgDckmO+BzRo= -github.com/hyperledger-firefly/transaction-manager v1.5.2/go.mod h1:OhqKmY1YWJhi2LZrGfl52ttRf3TfQYUZJ7j6nuSD06Q= +github.com/hyperledger-firefly/transaction-manager v1.5.3 h1:bnvwiVu6Jnj01MltTI9c7mqtwuWznVr33UDNbRH1oVU= +github.com/hyperledger-firefly/transaction-manager v1.5.3/go.mod h1:VlI9M+p6G4Z1btSyGIYSz/Jlz6LYfw8v35+g2jrRUSk= github.com/inconshreveable/mousetrap v1.1.0 h1:wN+x4NVGpMsO7ErUn/mUI3vEoE6Jt13X2s0bqwp9tc8= github.com/inconshreveable/mousetrap v1.1.0/go.mod h1:vpF70FUmC8bwa3OWnCshd2FqLfsEA9PFc4w1p2J65bw= github.com/jarcoal/httpmock v1.2.0 h1:gSvTxxFR/MEMfsGrvRbdfpRUMBStovlSRLw0Ep1bwwc= @@ -118,8 +118,8 @@ github.com/jarcoal/httpmock v1.2.0/go.mod h1:oCoTsnAz4+UoOUIf5lJOWV2QQIW5UoeUI6a github.com/karlseguin/ccache v2.0.3+incompatible h1:j68C9tWOROiOLWTS/kCGg9IcJG+ACqn5+0+t8Oh83UU= github.com/karlseguin/ccache v2.0.3+incompatible/go.mod h1:CM9tNPzT6EdRh14+jiW8mEF9mkNZuuE51qmgGYUB93w= github.com/kisielk/sqlstruct v0.0.0-20201105191214-5f3e10d3ab46/go.mod h1:yyMNCyc/Ib3bDTKd379tNMpB/7/H5TjM2Y9QJ5THLbE= -github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= -github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/klauspost/compress v1.19.1 h1:VsB4HPswih7mmZ8WleSFQ75c/Ui1M4trX5oAsJnhSlk= +github.com/klauspost/compress v1.19.1/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= @@ -160,8 +160,8 @@ github.com/oasdiff/yaml v0.1.1 h1:6nHx+pn9gBRM6YpBlFZFQGCCd1nuvqOBtTD3KKTgGxY= github.com/oasdiff/yaml v0.1.1/go.mod h1:EYJNoyktvWMJ0Hmhx+6qTaqMOsalUaRGT8Sj1hNcegU= github.com/oasdiff/yaml3 v0.0.14 h1:aLJee3hxBK2H5wdXd9iPcIXb93Nty1Ge0pT171eHtkw= github.com/oasdiff/yaml3 v0.0.14/go.mod h1:csto2xfDjYccdUn/yw/bPjj/cYTdp6HtFA0J4TWG+gg= -github.com/oklog/ulid/v2 v2.1.1 h1:suPZ4ARWLOJLegGFiZZ1dFAkqzhMjL3J1TzI+5wHz8s= -github.com/oklog/ulid/v2 v2.1.1/go.mod h1:rcEKHmBBKfef9DhnvX7y1HZBYxjXb0cP5ExxNsTT1QQ= +github.com/oklog/ulid/v2 v2.1.2 h1:IEclFb9JNvzYA6MW2SCxbLzcHTVsfqm3PrqGQJH5zec= +github.com/oklog/ulid/v2 v2.1.2/go.mod h1:rcEKHmBBKfef9DhnvX7y1HZBYxjXb0cP5ExxNsTT1QQ= github.com/onsi/ginkgo v1.6.0/go.mod h1:lLunBs/Ym6LB5Z9jYTR76FiuTmxDTDusOGeTQH+WWjE= github.com/onsi/ginkgo v1.12.1/go.mod h1:zj2OWP4+oCPe1qIXoGWkgMRwljMUYCdkwsT2108oapk= github.com/onsi/ginkgo v1.14.0/go.mod h1:iSB4RoI2tjJc9BBv4NKIKWKya62Rps+oPG/Lv9klQyY= @@ -183,12 +183,12 @@ github.com/pelletier/go-toml/v2 v2.2.4/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8 github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4= github.com/pkg/errors v0.9.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0= github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= -github.com/prometheus/client_golang v1.24.0 h1:5XStIklKuAtJSNpdD3s8XJj/Yv78IQmE1kbNk87JrAI= -github.com/prometheus/client_golang v1.24.0/go.mod h1:QcsNdotprC2nS4BTM2ucbcqxd2CeXTEa9jW7zHO9iDE= +github.com/prometheus/client_golang v1.24.1 h1:JnJkREXzWxUdCuPFpIWZiPispT9xVV59uiuyR2bPlnU= +github.com/prometheus/client_golang v1.24.1/go.mod h1:F+oSRECHg4sse5ucfYpYDeIv/hu68Zo0uoHKetWnzcE= github.com/prometheus/client_model v0.6.2 h1:oBsgwpGs7iVziMvrGhE53c/GrLUsZdHnqNwqPLxwZyk= github.com/prometheus/client_model v0.6.2/go.mod h1:y3m2F6Gdpfy6Ut/GBsUqTWZqCUvMVzSfMLjcu6wAwpE= -github.com/prometheus/common v0.70.0 h1:bcpru3tWPVnxGnETLgOV5jbp/JRXgYEyv65CuBLAMMI= -github.com/prometheus/common v0.70.0/go.mod h1:S/SFasQmgGiYH6C81LKCtYa8QACgthGg5zxL2udV7SY= +github.com/prometheus/common v0.70.1 h1:1HvjP4D5oL3t8RsPlwxA9onvvStjtIHYE5XuuwOi/PY= +github.com/prometheus/common v0.70.1/go.mod h1:VdFUQDMZK3VLkurFUVhia6uys/0suUp86TJz5qbJRhc= github.com/prometheus/procfs v0.21.1 h1:GljZCt+zSTS+NZq88cyQ1LjZ+RCHp3uVuabBWA5+OJI= github.com/prometheus/procfs v0.21.1/go.mod h1:aB55Cww9pdSJVHk0hUf0inxWyyjPogFIjmHKYgMKmtY= github.com/rogpeppe/go-internal v1.9.0 h1:73kH8U+JUqXU8lRuOHeVHaa/SZPifC7BkcraZVejAe8= @@ -204,8 +204,8 @@ github.com/santhosh-tekuri/jsonschema/v6 v6.0.2 h1:KRzFb2m7YtdldCEkzs6KqmJw4nqEV github.com/santhosh-tekuri/jsonschema/v6 v6.0.2/go.mod h1:JXeL+ps8p7/KNMjDQk3TCwPpBy0wYklyWTfbkIzdIFU= github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp81k= github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME= -github.com/sirupsen/logrus v1.10.0 h1:T8MxJJXVZkfcC5zSRMRAg2F8+lxjmUCGGWPzFxO+Msc= -github.com/sirupsen/logrus v1.10.0/go.mod h1:FXZFonkDAnFozmO+5hGAFvB0Yg9/j2SIhA/QuIkP180= +github.com/sirupsen/logrus v1.10.1 h1:xi4336Zh11WpU14fXR6I67V3yaTPQYwRx2WEtHbRg4Q= +github.com/sirupsen/logrus v1.10.1/go.mod h1:vsQHnG7xzNsxk3NrwboUiWPnIC3dmbjcGPykD7+tiHk= github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8 h1:+jumHNA0Wrelhe64i8F6HNlS8pkoyMv5sreGx2Ry5Rw= github.com/sourcegraph/conc v0.3.1-0.20240121214520-5f936abd7ae8/go.mod h1:3n1Cwaq1E1/1lhQhtRK2ts/ZwZEhjcQeJQ1RuC6Q/8U= github.com/spf13/afero v1.15.0 h1:b/YBCLWAJdFWJTN9cLhiXXcD7mzKn9Dm86dNnfyQw1I= diff --git a/internal/ethereum/event_actions_test.go b/internal/ethereum/event_actions_test.go index 4126b1a..7d971ae 100644 --- a/internal/ethereum/event_actions_test.go +++ b/internal/ethereum/event_actions_test.go @@ -159,6 +159,26 @@ func TestEventStreamStartStopOk(t *testing.T) { }) assert.NoError(t, err) assert.Equal(t, int64(testHighBlock), r2.Checkpoint.(*listenerCheckpoint).Block) + // Nothing has been pushed to FFTM for this listener, so we must report an untyped nil + assert.True(t, r2.LastDetected == nil) + + // Once we have pushed an event, we report its checkpoint as the detection point + c.eventStreams[*sID].listeners[*lID].markDetected(&listenerCheckpoint{ + Block: testHighBlock, + TransactionIndex: 123, + LogIndex: 0, + }) + r2b, _, err := c.EventListenerHWM(ctx, &ffcapi.EventListenerHWMRequest{ + StreamID: sID, + ListenerID: lID, + }) + assert.NoError(t, err) + assert.Equal(t, int64(testHighBlock), r2b.Checkpoint.(*listenerCheckpoint).Block) + assert.Equal(t, &listenerCheckpoint{ + Block: testHighBlock, + TransactionIndex: 123, + LogIndex: 0, + }, r2b.LastDetected) _, _, err = c.EventStreamStopped(ctx, &ffcapi.EventStreamStoppedRequest{ ID: sID, diff --git a/internal/ethereum/event_listener.go b/internal/ethereum/event_listener.go index 25d9058..6fdcd20 100644 --- a/internal/ethereum/event_listener.go +++ b/internal/ethereum/event_listener.go @@ -60,8 +60,9 @@ type listener struct { c *ethConnector es *eventStream ee *eventEnricher - hwmMux sync.Mutex // Protects checkpoint of an individual listener. May hold ES lock when taking this, must NOT attempt to obtain ES lock while holding this - hwmBlock int64 + hwmMux sync.Mutex // Protects hwmBlock and lastDetected. May hold ES lock when taking this, must NOT attempt to obtain ES lock while holding this + hwmBlock int64 // the scan position - the block we have polled for events up to (exclusive) + lastDetected *listenerCheckpoint // the checkpoint of the highest event we have pushed to FFTM - see ffcapi.EventListenerHWMResponse config listenerConfig removed bool catchup bool @@ -130,20 +131,35 @@ func (l *listener) checkReadyForLeadPackOrRemoved(ctx context.Context) (bool, bo return readyForLead, l.removed } -// getHWMCheckpoint gets the point the event polling is up to for this listener. -// Note this intentionally does not account for dispatched events, as the parent framework ensures that -// this checkpoint is only persisted when there are no events in-flight pending dispatch for this listener, -// and the checkpoint for this listener is stale. -func (l *listener) getHWMCheckpoint() *listenerCheckpoint { +// getHWM returns under the hmwMux lock as consistent set of: +// 1. Scan position - where event polling is up to, for detection of new events +// 2. Last detected - the checkpoint of the highest event pushed to the FFTM channel (or nil) +// See ffcapi.EventListenerHWMResponse for the contract defined by FFTM for these +func (l *listener) getHWM() (scanned ffcapi.EventListenerCheckpoint, lastDetected ffcapi.EventListenerCheckpoint) { l.hwmMux.Lock() defer l.hwmMux.Unlock() // Generate a checkpoint before the first transaction, in the high watermark block - log.L(l.es.ctx).Debugf("HWM checkpoint block for '%s': %d", l.id, l.hwmBlock) - return &listenerCheckpoint{ + log.L(l.es.ctx).Debugf("HWM checkpoint block for '%s': %d (lastDetected=%+v)", l.id, l.hwmBlock, l.lastDetected) + scanned = &listenerCheckpoint{ Block: l.hwmBlock, TransactionIndex: -1, LogIndex: -1, } + if l.lastDetected != nil { + lastDetected = l.lastDetected + } + return scanned, lastDetected +} + +// markDetected records the checkpoint of an event, and must be called pushing the event to FFTM +func (l *listener) markDetected(cp *listenerCheckpoint) { + l.hwmMux.Lock() + defer l.hwmMux.Unlock() + // Only ever move forwards - a re-detection (such as after a re-org, or a filter reset) must not + // lower the bar FFTM uses to decide the scan position is safe to record as a checkpoint. + if l.lastDetected == nil || l.lastDetected.LessThan(cp) { + l.lastDetected = cp + } } func (l *listener) moveHWMForwards(hwmBlock int64) { @@ -231,10 +247,7 @@ func (l *listener) listenerCatchupLoop() { log.L(ctx).Infof("Listener catchup fromBlock=%d toBlock=%d events=%d", fromBlock, toBlock, len(events)) for _, event := range events { - log.L(ctx).Debugf("Detected event %s (listener catchup)", event.Event) - select { - case l.es.events <- event: - case <-l.es.ctx.Done(): + if l.es.markDetectedAndDispatch(al, event) { log.L(ctx).Infof("Listener catchup loop exiting as stream is stopping") return } diff --git a/internal/ethereum/event_listener_test.go b/internal/ethereum/event_listener_test.go index 53cd327..f17de22 100644 --- a/internal/ethereum/event_listener_test.go +++ b/internal/ethereum/event_listener_test.go @@ -681,3 +681,40 @@ func TestFilterEnrichEthLogMethodBadInputABIData(t *testing.T) { assert.Nil(t, ei.InputArgs) } + +func TestGetHWMNoDetectionsYet(t *testing.T) { + + l, _, cancelCtx := newTestListener(t, false) + defer cancelCtx() + + l.hwmBlock = 1000 + + scanned, lastDetected := l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: -1, LogIndex: -1}, scanned) + // Must be an untyped nil interface - a nil *listenerCheckpoint inside a non-nil interface + // would make FFTM believe we had detected something, and panic comparing against it + assert.True(t, lastDetected == nil) + +} + +func TestMarkDetectedOnlyMovesForwards(t *testing.T) { + + l, _, cancelCtx := newTestListener(t, false) + defer cancelCtx() + + l.markDetected(&listenerCheckpoint{Block: 1000, TransactionIndex: 5, LogIndex: 0}) + _, lastDetected := l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: 5, LogIndex: 0}, lastDetected) + + // A later event in the same block moves it forwards + l.markDetected(&listenerCheckpoint{Block: 1000, TransactionIndex: 7, LogIndex: 0}) + _, lastDetected = l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: 7, LogIndex: 0}, lastDetected) + + // A re-detection of an earlier event does not lower the bar FFTM uses to decide + // the scan position is safe to apply + l.markDetected(&listenerCheckpoint{Block: 999, TransactionIndex: 0, LogIndex: 0}) + _, lastDetected = l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: 7, LogIndex: 0}, lastDetected) + +} diff --git a/internal/ethereum/event_stream.go b/internal/ethereum/event_stream.go index 081cb65..e349318 100644 --- a/internal/ethereum/event_stream.go +++ b/internal/ethereum/event_stream.go @@ -75,6 +75,7 @@ type eventStream struct { type aggregatedListener struct { signatureSet []ethtypes.HexBytes0xPrefix // a list of unique topic[0] event signatures to listener for listenersByTopic0 map[string][]*listener // a map of all listeners that are interested in an event signature - they may not be interested in the event itself (depending on sub-selection) + listenersByID map[fftypes.UUID]*listener // a map of all listeners by ID, to resolve the listener that generated an event when dispatching it listeners []*listener // list of all listeners } @@ -506,10 +507,7 @@ func (es *eventStream) dispatchSetHWMCheckExit(ag *aggregatedListener, events ff } } else { for _, event := range events { - log.L(es.ctx).Debugf("Detected event %s", event.Event) - select { - case es.events <- event: - case <-es.ctx.Done(): + if es.markDetectedAndDispatch(ag, event) { return true } } @@ -524,12 +522,29 @@ func (es *eventStream) dispatchSetHWMCheckExit(ag *aggregatedListener, events ff } +// markDetectedAndDispatch records the detection point then (importantly afterwards) pushes the event to FFTM +func (es *eventStream) markDetectedAndDispatch(ag *aggregatedListener, event *ffcapi.ListenerEvent) (exiting bool) { + log.L(es.ctx).Debugf("Detected event %s", event.Event) + + // ListenerID is set in filterEnrichEthLog and must be non-nil + ag.listenersByID[*event.Event.ID.ListenerID].markDetected(event.Checkpoint.(*listenerCheckpoint)) + select { + case es.events <- event: + return false + case <-es.ctx.Done(): + return true + } + +} + func (es *eventStream) buildAggregatedListener(listeners []*listener) *aggregatedListener { ag := &aggregatedListener{ listeners: listeners, listenersByTopic0: make(map[string][]*listener), + listenersByID: make(map[fftypes.UUID]*listener), } for _, l := range listeners { + ag.listenersByID[*l.id] = l for _, f := range l.config.filters { sigStr := f.Topic0.String() topicListeners, existing := ag.listenersByTopic0[sigStr] @@ -595,8 +610,10 @@ func (es *eventStream) getListenerHWM(ctx context.Context, listenerID *fftypes.U if l == nil { return nil, ffcapi.ErrorReasonNotFound, i18n.NewError(ctx, msgs.MsgListenerNotStarted, listenerID, es.id) } + scanned, lastDetected := l.getHWM() return &ffcapi.EventListenerHWMResponse{ - Checkpoint: l.getHWMCheckpoint(), - Catchup: l.catchup || es.catchup, // dirty read of whether the listener is in catchup, or the head group of the stream is in catchup + Checkpoint: scanned, + LastDetected: lastDetected, + Catchup: l.catchup || es.catchup, // dirty read of whether the listener is in catchup, or the head group of the stream is in catchup }, "", nil } diff --git a/internal/ethereum/event_stream_test.go b/internal/ethereum/event_stream_test.go index 57bed37..cc5aa74 100644 --- a/internal/ethereum/event_stream_test.go +++ b/internal/ethereum/event_stream_test.go @@ -1221,11 +1221,21 @@ func TestDispatchListenerDone(t *testing.T) { ctx: doneCtx, events: make(chan<- *ffcapi.ListenerEvent), } - exiting := es.dispatchSetHWMCheckExit(&aggregatedListener{}, ffcapi.ListenerEvents{ - {}, + l := &listener{id: fftypes.NewUUID(), es: es} + ag := es.buildAggregatedListener([]*listener{l}) + exiting := es.dispatchSetHWMCheckExit(ag, ffcapi.ListenerEvents{ + { + Checkpoint: &listenerCheckpoint{Block: 1000, TransactionIndex: 10, LogIndex: 1}, + Event: &ffcapi.Event{ID: ffcapi.EventID{ListenerID: l.id}}, + }, }, -1) assert.True(t, exiting) + // The detection point is recorded before the dispatch, so it survives losing the ctx race - + // this can only ever hold the checkpoint back, never advance it past an undelivered event + _, lastDetected := l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: 10, LogIndex: 1}, lastDetected) + } func TestGetListenerHWMNotFound(t *testing.T) { @@ -1240,3 +1250,35 @@ func TestGetListenerHWMNotFound(t *testing.T) { assert.Equal(t, ffcapi.ErrorReasonNotFound, rc) } + +func TestDispatchSetHWMDetectionBeforeScanPosition(t *testing.T) { + + delivered := make(chan *ffcapi.ListenerEvent, 2) + es := &eventStream{ + ctx: context.Background(), + events: delivered, + } + l := &listener{id: fftypes.NewUUID(), es: es} + ag := es.buildAggregatedListener([]*listener{l}) + + exiting := es.dispatchSetHWMCheckExit(ag, ffcapi.ListenerEvents{ + { + Checkpoint: &listenerCheckpoint{Block: 1000, TransactionIndex: 5, LogIndex: 0}, + Event: &ffcapi.Event{ID: ffcapi.EventID{ListenerID: l.id}}, + }, + { + Checkpoint: &listenerCheckpoint{Block: 1000, TransactionIndex: 7, LogIndex: 0}, + Event: &ffcapi.Event{ID: ffcapi.EventID{ListenerID: l.id}}, + }, + }, 1001) + assert.False(t, exiting) + assert.Len(t, delivered, 2) + + // The scan position has moved past the block, but the detection point pins it to the highest + // event we pushed - so FFTM holds the checkpoint back until it has committed that far + scanned, lastDetected := l.getHWM() + assert.Equal(t, &listenerCheckpoint{Block: 1001, TransactionIndex: -1, LogIndex: -1}, scanned) + assert.Equal(t, &listenerCheckpoint{Block: 1000, TransactionIndex: 7, LogIndex: 0}, lastDetected) + assert.True(t, lastDetected.LessThan(scanned)) + +} diff --git a/pkg/ethblocklistener/block_receipt_fetcher_test.go b/pkg/ethblocklistener/block_receipt_fetcher_test.go index fd78c04..5a1b16b 100644 --- a/pkg/ethblocklistener/block_receipt_fetcher_test.go +++ b/pkg/ethblocklistener/block_receipt_fetcher_test.go @@ -28,6 +28,7 @@ import ( "github.com/hyperledger-firefly/signer/pkg/rpcbackend" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/mock" + "github.com/stretchr/testify/require" ) func TestFetchBlockReceiptsAsyncOptimizedOk(t *testing.T) { @@ -286,6 +287,34 @@ func TestReconcileConfirmationsForTransactionUsesCachedReceipt(t *testing.T) { mRPC.AssertNotCalled(t, "CallRPC", mock.Anything, mock.Anything, "eth_getTransactionReceipt", mock.Anything) } +func TestReceiptCacheDisabledAllEntryPointsNoop(t *testing.T) { + _, bl, _, done := newTestBlockListener(t, func(conf *BlockListenerConfig, mRPC *rpcbackendmocks.Backend, cancelCtx context.CancelFunc) { + conf.ReceiptCacheEnabled = false + }) + defer done() + + require.Nil(t, bl.txReceiptCache) + bl.canonicalChain = createTestChain(1976, 1978) + + txHash := "0x6197ef1a58a2a592bb447efb651f0db7945de21aa8048801b250bd7b7431f9b6" + bl.storeReceiptsInCache([]*ethrpc.TxReceiptJSONRPC{ + {TransactionHash: ethtypes.MustNewHexBytes0xPrefix(txHash)}, + }, bl.getReceiptCacheGeneration()) + + // Nothing is stored, so nothing can be served from the cache + _, ok := bl.getCachedTransactionReceipt(txHash) + assert.False(t, ok) + + // The reset and the two pre-fetch paths are all no-ops - the latter asserted by done(), + // as no receipt queries are mocked for the chain we put in place above + bl.resetReceiptCache() + bl.fetchAndCacheBlockReceipts(ðrpc.BlockInfoJSONRPC{ + Number: ethtypes.HexUint64(1977), + Hash: generateTestHash(1977), + }) + bl.refetchReceiptsForCanonicalChain() +} + func TestReceiptCacheEvictsWhenFull(t *testing.T) { _, bl, _, done := newTestBlockListener(t, func(conf *BlockListenerConfig, mRPC *rpcbackendmocks.Backend, cancelCtx context.CancelFunc) { conf.ReceiptCacheEnabled = true diff --git a/pkg/ethblocklistener/blocklistener_metrics_test.go b/pkg/ethblocklistener/blocklistener_metrics_test.go index a8e29f8..836e81d 100644 --- a/pkg/ethblocklistener/blocklistener_metrics_test.go +++ b/pkg/ethblocklistener/blocklistener_metrics_test.go @@ -158,11 +158,10 @@ func TestBlockListenerMetricsFullMode(t *testing.T) { ctx, bl, _, done := newTestBlockListener(t, func(conf *BlockListenerConfig, mRPC *rpcbackendmocks.Backend, _ context.CancelFunc) { conf.BlockPollingInterval = 1 * time.Millisecond - mRPC.On("CallRPC", mock.Anything, mock.Anything, "eth_blockNumber").Return(nil).Run(func(args mock.Arguments) { - *args[1].(*ethtypes.HexInteger) = *ethtypes.NewHexIntegerU64(1001) + *args[1].(*ethtypes.HexInteger) = *ethtypes.NewHexIntegerU64(1000) }) - mockSeedBlockNotFound(mRPC, 1001-uint64(conf.MonitoredHeadLength)+1) + mockSeedBlockNotFound(mRPC, 1000-uint64(conf.MonitoredHeadLength)+1) mockNewBlockFilter(mRPC, testBlockFilterID1) mockFilterChanges(mRPC, testBlockFilterID1, nil, blockHash1001).Once() mockFilterChangesEmpty(mRPC) @@ -180,7 +179,7 @@ func TestBlockListenerMetricsFullMode(t *testing.T) { }) // The height the node reports, refreshed by the listen loop, and the head of the chain we've built - waitForGaugeMetric(t, registry, metricTargetBlockHeight, 1001) + waitForGaugeMetric(t, registry, metricTargetBlockHeight, 1000) waitForGaugeMetric(t, registry, metricCanonicalBlockHeight, 1001) } diff --git a/pkg/ethblocklistener/blocklistener_test.go b/pkg/ethblocklistener/blocklistener_test.go index affffda..672cfc3 100644 --- a/pkg/ethblocklistener/blocklistener_test.go +++ b/pkg/ethblocklistener/blocklistener_test.go @@ -124,11 +124,15 @@ func mockInitialBlockHeight(mRPC *rpcbackendmocks.Backend, height uint64) *mock. }).Once() } -// mockSeedBlockNotFound mocks eth_getBlockByNumber at seedHeight returning nil (block not found). +// mockSeedBlockNotFound mocks one eth_getBlockByNumber at seedHeight returning nil (block not found). func mockSeedBlockNotFound(mRPC *rpcbackendmocks.Backend, seedHeight uint64) *mock.Call { return mRPC.On("CallRPC", mock.Anything, mock.Anything, "eth_getBlockByNumber", hexNumber(seedHeight), false).Return(nil).Once() } +func mockSeedBlockNotFoundMaybe(mRPC *rpcbackendmocks.Backend, seedHeight uint64) *mock.Call { + return mRPC.On("CallRPC", mock.Anything, mock.Anything, "eth_getBlockByNumber", hexNumber(seedHeight), false).Return(nil).Maybe() +} + // mockSeedBlock mocks eth_getBlockByNumber at height returning a block with the given hash. func mockSeedBlock(mRPC *rpcbackendmocks.Backend, height uint64, hash ethtypes.HexBytes0xPrefix) *mock.Call { return mRPC.On("CallRPC", mock.Anything, mock.Anything, "eth_getBlockByNumber", hexNumber(height), false).Return(nil).Run(func(args mock.Arguments) { @@ -250,7 +254,7 @@ func TestBlockListenerStartGettingHighestBlockRetry(t *testing.T) { mRPC.On("CallRPC", mock.Anything, mock.Anything, "eth_blockNumber"). Return(&rpcbackend.RPCError{Message: "pop"}).Once() mockInitialBlockHeight(mRPC, 12345) - mockSeedBlockNotFound(mRPC, 12345-(50-1)).Maybe() + mockSeedBlockNotFoundMaybe(mRPC, 12345-(50-1)) mockNewBlockFilter(mRPC, testBlockFilterID1).Maybe() mockFilterChangesEmpty(mRPC).Maybe() @@ -1397,7 +1401,7 @@ func TestGetHighestBlockInfoBeforeHeadBlockSeen(t *testing.T) { _, bl, mRPC, done := newTestBlockListener(t) mockInitialBlockHeight(mRPC, 500) - mockSeedBlockNotFound(mRPC, 500-(50-1)).Maybe() + mockSeedBlockNotFoundMaybe(mRPC, 500-(50-1)) mockNewBlockFilter(mRPC, testBlockFilterID1).Maybe() mockFilterChangesEmpty(mRPC).Maybe() @@ -1415,7 +1419,7 @@ func TestGetHighestBlockInfoReturnsHeadBlock(t *testing.T) { _, bl, mRPC, done := newTestBlockListener(t) mockInitialBlockHeight(mRPC, 123) - mockSeedBlockNotFound(mRPC, 123-(50-1)).Maybe() + mockSeedBlockNotFoundMaybe(mRPC, 123-(50-1)) mockNewBlockFilter(mRPC, testBlockFilterID1).Maybe() mockFilterChangesEmpty(mRPC).Maybe()