Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions config.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`|`<nil>`
Expand Down
12 changes: 6 additions & 6 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -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
Expand Down
28 changes: 14 additions & 14 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down Expand Up @@ -109,17 +109,17 @@ 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=
github.com/jarcoal/httpmock v1.2.0/go.mod h1:oCoTsnAz4+UoOUIf5lJOWV2QQIW5UoeUI6aM2YnWAZk=
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=
Expand Down Expand Up @@ -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=
Expand All @@ -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=
Expand All @@ -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=
Expand Down
20 changes: 20 additions & 0 deletions internal/ethereum/event_actions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
39 changes: 26 additions & 13 deletions internal/ethereum/event_listener.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -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
}
Expand Down
37 changes: 37 additions & 0 deletions internal/ethereum/event_listener_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)

}
29 changes: 23 additions & 6 deletions internal/ethereum/event_stream.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
}

Expand Down Expand Up @@ -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
}
}
Expand All @@ -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))
Comment thread
peterbroadhurst marked this conversation as resolved.
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]
Expand Down Expand Up @@ -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
}
46 changes: 44 additions & 2 deletions internal/ethereum/event_stream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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))

}
Loading