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
2 changes: 1 addition & 1 deletion Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ $(eval $(call makemock, internal/persistence, TransactionPersistence, per
$(eval $(call makemock, internal/persistence, RichQuery, persistencemocks))
$(eval $(call makemock, internal/ws, WebSocketChannels, wsmocks))
$(eval $(call makemock, internal/ws, WebSocketServer, wsmocks))
$(eval $(call makemock, internal/events, Stream, eventsmocks))
$(eval $(call makemock, internal/events, ManagedStream, eventsmocks))
$(eval $(call makemock, internal/apiclient, FFTMClient, apiclientmocks))

go-mod-tidy: .ALWAYS
Expand Down
123 changes: 91 additions & 32 deletions internal/events/eventstream.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,29 @@ import (

type Stream = eventapi.EventStream

// ManagedStream extends the connector-facing stream interface with the extra functions FFTM's
// manager needs - which cannot be added to pkg/eventapi without breaking connectors.
type ManagedStream interface {
Stream

// VerifyListenerOptions is used in create/update to verify and resolve the full spec
// including the Name, which defaults to the resolved signature
VerifyListenerOptions(ctx context.Context, id *fftypes.UUID, updatesOrNew *apitypes.Listener) (*apitypes.Listener, error)

// PrepareListenerUpdate checks a spec previously returned by VerifyListenerOptions against the
// current state of the stream, without changing anything. The caller persists the returned
// Spec, then calls Apply() to make the update live.
PrepareListenerUpdate(ctx context.Context, spec *apitypes.Listener, reset bool) (*PreparedListenerUpdate, error)
}

// PreparedListenerUpdate is a listener create/update that has been validated against a stream, but
// not yet applied. This allows the change to be persisted to the database after prepare, before apply.
type PreparedListenerUpdate struct {
Spec *apitypes.Listener // new/updated spec fully validated
es *eventStream
reset bool
}

// esDefaults are the defaults for new event streams, read from the config once in InitDefaults()
var esDefaults struct {
initialized bool
Expand Down Expand Up @@ -123,7 +146,7 @@ func NewEventStream(
wsChannels ws.WebSocketChannels,
initialListeners []*apitypes.Listener,
eme metrics.EventMetricsEmitter,
) (ees Stream, err error) {
) (ees ManagedStream, err error) {
return newEventStream(
bgCtx,
persistedSpec,
Expand All @@ -141,7 +164,7 @@ func NewAPIManagedEventStream(
connector ffcapi.API,
listeners []*apitypes.Listener,
eme metrics.EventMetricsEmitter,
) (ees Stream, err error) {
) (ees ManagedStream, err error) {
return newEventStream(
bgCtx,
persistedSpec,
Expand All @@ -162,7 +185,7 @@ func newEventStream(
wsChannels ws.WebSocketChannels,
initialListeners []*apitypes.Listener,
eme metrics.EventMetricsEmitter,
) (ees Stream, err error) {
) (ees ManagedStream, err error) {
esCtx := log.WithLogFields(bgCtx, "eventstream", spec.ID.String())
es := &eventStream{
bgCtx: esCtx,
Expand All @@ -185,7 +208,7 @@ func newEventStream(
}
es.batchChannel = make(chan *ffcapi.ListenerEvent, *es.spec.BatchSize)
for _, existing := range initialListeners {
spec, err := es.verifyListenerOptions(esCtx, existing.ID, existing)
spec, err := es.VerifyListenerOptions(esCtx, existing.ID, existing)
if err != nil {
return nil, err
}
Expand Down Expand Up @@ -395,7 +418,7 @@ func (es *eventStream) mergeListenerOptions(id *fftypes.UUID, updates *apitypes.

}

func (es *eventStream) verifyListenerOptions(ctx context.Context, id *fftypes.UUID, updatesOrNew *apitypes.Listener) (*apitypes.Listener, error) {
func (es *eventStream) VerifyListenerOptions(ctx context.Context, id *fftypes.UUID, updatesOrNew *apitypes.Listener) (*apitypes.Listener, error) {
if id == nil {
return nil, i18n.NewError(ctx, tmmsgs.MsgMissingID)
}
Expand Down Expand Up @@ -437,49 +460,91 @@ func (es *eventStream) verifyListenerOptions(ctx context.Context, id *fftypes.UU
return spec, nil
}

// AddOrUpdateListener resolves and applies a listener in one step.
// Note this is only used externally.
// The FFTM manager needs to coordinate DB work to persist the change
// between VerifyListenerOptions and PrepareListenerUpdate.
func (es *eventStream) AddOrUpdateListener(ctx context.Context, id *fftypes.UUID, updates *apitypes.Listener, reset bool) (merged *apitypes.Listener, err error) {
log.L(ctx).Infof("Adding/updating listener %s", id)

// Ask the connector to verify the options, and apply defaults
spec, err := es.verifyListenerOptions(ctx, id, updates)
spec, err := es.VerifyListenerOptions(ctx, id, updates)
if err != nil {
return nil, err
}

// Do the locked part - which checks if this is a new listener, or just an update to the options.
isNew, l, startedState, err := es.lockedListenerUpdate(ctx, spec, reset)
u, err := es.PrepareListenerUpdate(ctx, spec, reset)
if err != nil {
return nil, err
}
if reset {
if err := u.Apply(ctx); err != nil {
return nil, err
}
return u.Spec, nil
}

// PrepareListenerUpdate checks that a resolved spec is, and returns an update ready to Apply().
func (es *eventStream) PrepareListenerUpdate(ctx context.Context, spec *apitypes.Listener, reset bool) (*PreparedListenerUpdate, error) {
log.L(ctx).Infof("Preparing listener %s", spec.ID)

es.mux.Lock()
defer es.mux.Unlock()

l, exists := es.listeners[*spec.ID]
switch {
case exists:
if spec.SignatureString() != l.spec.SignatureString() {
// We do not allow the filters to be updated, because that would lead to a confusing situation
// where the previously emitted events are a subset/mismatch to the filters configured now.
return nil, i18n.NewError(ctx, tmmsgs.MsgFilterUpdateNotAllowed, l.spec.SignatureString(), spec.SignatureString())
}
case reset:
return nil, i18n.NewError(ctx, tmmsgs.MsgResetStreamNotFound, spec.ID, es.spec.ID)
}

return &PreparedListenerUpdate{
Spec: spec,
es: es,
reset: reset,
}, nil
}

// Apply finalizes the update so this new version is the one in the in-memory map of listeners
// for the stream, and resetting the checkpoint if required.
func (u *PreparedListenerUpdate) Apply(ctx context.Context) error {
es := u.es

// Update the map to apply this version of the spec (new or updated based on the state at apply)
l, isNew, startedState := es.lockedListenerUpdate(u.Spec)
log.L(ctx).Infof("Applied listener %s (new=%t, reset=%t)", u.Spec.ID, isNew, u.reset)

switch {
case u.reset:
// Only safe to do the reset with the event stream stopped
if startedState != nil {
if err := es.Stop(ctx); err != nil {
return nil, err
return err
}
}
// Clear out the checkpoint for this listener
if err := es.resetListenerCheckpoint(ctx, l); err != nil {
return nil, err
return err
}
// Restart if we were started
if startedState != nil {
if err := es.Start(ctx); err != nil {
return nil, err
return err
}
}
} else if isNew && startedState != nil {
case isNew && startedState != nil:
if l.spec.Type != nil && *l.spec.Type == apitypes.ListenerTypeBlocks {
err := l.es.confirmations.StartConfirmedBlockListener(ctx, l.spec.ID, *l.spec.FromBlock, nil /* new so no checkpoint */, es.batchChannel)
if err == nil {
l.markStarted(true)
}
return spec, err
return err
}
// Start the new listener - no checkpoint needed here
return spec, l.start(startedState, nil)
return l.start(startedState, nil)
}
return spec, nil
return nil
}

func (es *eventStream) resetListenerCheckpoint(ctx context.Context, l *listener) error {
Expand All @@ -493,30 +558,24 @@ func (es *eventStream) resetListenerCheckpoint(ctx context.Context, l *listener)
return es.checkpointsDB.WriteCheckpoint(ctx, cp)
}

func (es *eventStream) lockedListenerUpdate(ctx context.Context, spec *apitypes.Listener, reset bool) (bool, *listener, *startedStreamState, error) {
// lockedListenerUpdate adds the listener to the stream, or swaps the spec of the existing one,
// returning it along with the started state read in the same lock hold - so the caller starts it
// against the state it was actually added to.
func (es *eventStream) lockedListenerUpdate(spec *apitypes.Listener) (l *listener, isNew bool, startedState *startedStreamState) {
es.mux.Lock()
defer es.mux.Unlock()

l, exists := es.listeners[*spec.ID]
switch {
case exists:
if spec.SignatureString() != l.spec.SignatureString() {
// We do not allow the filters to be updated, because that would lead to a confusing situation
// where the previously emitted events are a subset/mismatch to the filters configured now.
return false, nil, nil, i18n.NewError(ctx, tmmsgs.MsgFilterUpdateNotAllowed, l.spec.SignatureString(), spec.SignatureString())
}
if exists {
l.spec = spec
case reset:
return false, nil, nil, i18n.NewError(ctx, tmmsgs.MsgResetStreamNotFound, spec.ID, es.spec.ID)
default:
} else {
l = &listener{
es: es,
spec: spec,
}
es.listeners[*spec.ID] = l
}
// Take a copy of the current started status, before unlocking
return !exists, l, es.currentState, nil
return l, !exists, es.currentState
}

func (es *eventStream) RemoveListener(ctx context.Context, id *fftypes.UUID) (err error) {
Expand Down
145 changes: 145 additions & 0 deletions internal/events/eventstream_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2468,3 +2468,148 @@ func TestStartAPIEventStreamPollContextCancelled(t *testing.T) {
require.Regexp(t, "FF00154", err)

}

// Check in the two-stage verify+apply, the original listener stays in-place until the Apply() completes
func TestPrepareListenerUpdateInvisibleUntilApply(t *testing.T) {

es := newTestEventStream(t, `{
"name": "ut_stream"
}`)

mfc := es.connector.(*ffcapimocks.API)
mfc.On("EventListenerVerifyOptions", mock.Anything, mock.Anything).Return(&ffcapi.EventListenerVerifyOptionsResponse{
ResolvedSignature: "sig1",
}, ffcapi.ErrorReason(""), nil)
initialListeners := make(chan int, 10)
mfc.On("EventStreamStart", mock.Anything, mock.Anything).Run(func(args mock.Arguments) {
initialListeners <- len(args[1].(*ffcapi.EventStreamStartRequest).InitialListeners)
}).Return(&ffcapi.EventStreamStartResponse{}, ffcapi.ErrorReason(""), nil)
added := make(chan *fftypes.UUID, 1)
mfc.On("EventListenerAdd", mock.Anything, mock.Anything).Run(func(args mock.Arguments) {
added <- args[1].(*ffcapi.EventListenerAddRequest).ListenerID
}).Return(nil, ffcapi.ErrorReason(""), nil)
mfc.On("EventStreamStopped", mock.Anything, mock.Anything).Return(&ffcapi.EventStreamStoppedResponse{}, ffcapi.ErrorReason(""), nil).Maybe()

msp := es.checkpointsDB.(*persistencemocks.Persistence)
msp.On("GetCheckpoint", mock.Anything, mock.Anything).Return(nil, nil)

require.NoError(t, es.Start(es.bgCtx))
assert.Equal(t, 0, <-initialListeners)

lID := fftypes.NewUUID()
spec, err := es.VerifyListenerOptions(es.bgCtx, lID, &apitypes.Listener{
Name: strPtr("ut_listener"),
Filters: []fftypes.JSONAny{`{"event":"definition1"}`},
})
require.NoError(t, err)
u, err := es.PrepareListenerUpdate(es.bgCtx, spec, false)
require.NoError(t, err)

// The stream does not know about the listener yet
es.mux.Lock()
_, exists := es.listeners[*lID]
es.mux.Unlock()
assert.False(t, exists)

// ... so a restart in the window cannot register it with the connector
require.NoError(t, es.Stop(es.bgCtx))
require.NoError(t, es.Start(es.bgCtx))
assert.Equal(t, 0, <-initialListeners)

// Applying is what makes it real
require.NoError(t, u.Apply(es.bgCtx))
assert.Equal(t, lID, <-added)
es.mux.Lock()
l, exists := es.listeners[*lID]
assert.True(t, l.started)
es.mux.Unlock()
require.True(t, exists)

require.NoError(t, es.Stop(es.bgCtx))
}

func TestPrepareListenerUpdateExistingUnchangedUntilApply(t *testing.T) {

es := newTestEventStream(t, `{
"name": "ut_stream"
}`)

mfc := es.connector.(*ffcapimocks.API)
mfc.On("EventListenerVerifyOptions", mock.Anything, mock.Anything).Return(&ffcapi.EventListenerVerifyOptionsResponse{
ResolvedSignature: "sig1",
}, ffcapi.ErrorReason(""), nil)
// An update must never remove the listener it is updating
mfc.On("EventListenerRemove", mock.Anything, mock.Anything).Run(func(args mock.Arguments) {
assert.Fail(t, "listener must not be removed by an update")
}).Return(&ffcapi.EventListenerRemoveResponse{}, ffcapi.ErrorReason(""), nil).Maybe()

lID := fftypes.NewUUID()
original, err := es.VerifyListenerOptions(es.bgCtx, lID, &apitypes.Listener{
Name: strPtr("name1"),
Filters: []fftypes.JSONAny{`{"event":"definition1"}`},
})
require.NoError(t, err)
u, err := es.PrepareListenerUpdate(es.bgCtx, original, false)
require.NoError(t, err)
require.NoError(t, u.Apply(es.bgCtx))

updated, err := es.VerifyListenerOptions(es.bgCtx, lID, &apitypes.Listener{
Name: strPtr("name2"),
})
require.NoError(t, err)
u, err = es.PrepareListenerUpdate(es.bgCtx, updated, false)
require.NoError(t, err)

// The listener is still on its previous spec until the caller applies
es.mux.Lock()
l := es.listeners[*lID]
es.mux.Unlock()
assert.Same(t, original, l.spec)

require.NoError(t, u.Apply(es.bgCtx))
es.mux.Lock()
l, exists := es.listeners[*lID]
es.mux.Unlock()
require.True(t, exists)
assert.Same(t, updated, l.spec)
}

func TestApplyListenerUpdateRemovedWhilePersisting(t *testing.T) {

es := newTestEventStream(t, `{
"name": "ut_stream"
}`)

mfc := es.connector.(*ffcapimocks.API)
mfc.On("EventListenerVerifyOptions", mock.Anything, mock.Anything).Return(&ffcapi.EventListenerVerifyOptionsResponse{
ResolvedSignature: "sig1",
}, ffcapi.ErrorReason(""), nil)

msp := es.checkpointsDB.(*persistencemocks.Persistence)
msp.On("GetCheckpoint", mock.Anything, mock.Anything).Return(nil, nil)

lID := fftypes.NewUUID()
spec, err := es.VerifyListenerOptions(es.bgCtx, lID, &apitypes.Listener{
Name: strPtr("name1"),
Filters: []fftypes.JSONAny{`{"event":"definition1"}`},
})
require.NoError(t, err)
u, err := es.PrepareListenerUpdate(es.bgCtx, spec, false)
require.NoError(t, err)
require.NoError(t, u.Apply(es.bgCtx))

// A reset prepared against a listener that is removed before we apply. The caller has written
// the spec by this point, so the right answer is to end up matching it, not to fail.
u, err = es.PrepareListenerUpdate(es.bgCtx, spec, true)
require.NoError(t, err)
es.mux.Lock()
delete(es.listeners, *lID)
es.mux.Unlock()
require.NoError(t, u.Apply(es.bgCtx))

es.mux.Lock()
l, exists := es.listeners[*lID]
es.mux.Unlock()
require.True(t, exists)
assert.Same(t, spec, l.spec)
}
1 change: 1 addition & 0 deletions internal/tmmsgs/en_error_messages.go
Original file line number Diff line number Diff line change
Expand Up @@ -112,4 +112,5 @@ var (
MsgConfirmedBlockListenerUnsupportedMode = ffe("FF21095", "Block listener is not supported when chain tracking mode is '%s'", http.StatusBadRequest)
MsgTransactionReceiptMissingBlockHash = ffe("FF21096", "Transaction receipt missing block hash")
MsgTransactionReceiptBlockHashMismatch = ffe("FF21097", "Transaction receipt block hash mismatch. Expected %s, got %s")
MsgDuplicateListenerName = ffe("FF21098", "Duplicate listener name '%s' used by listener '%s'", http.StatusConflict)
)
Loading