mirror of
https://github.com/tendermint/tendermint.git
synced 2026-09-29 11:15:50 +00:00
types: rename and extend the EventData interface (#7687)
This is the interface shared by types that can be used as event data in, for example, subscriptions via the RPC. To be compatible with the RPC service, data need to support JSON encoding. Require this as part of the interface.
This commit is contained in:
@@ -12,12 +12,19 @@ import (
|
||||
"github.com/tendermint/tendermint/internal/pubsub"
|
||||
"github.com/tendermint/tendermint/internal/pubsub/query"
|
||||
"github.com/tendermint/tendermint/libs/log"
|
||||
"github.com/tendermint/tendermint/types"
|
||||
)
|
||||
|
||||
const (
|
||||
clientID = "test-client"
|
||||
)
|
||||
|
||||
// pubstring is a trivial implementation of the EventData interface for
|
||||
// string-valued test data.
|
||||
type pubstring string
|
||||
|
||||
func (pubstring) TypeTag() string { return "pubstring" }
|
||||
|
||||
func TestSubscribeWithArgs(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
@@ -34,8 +41,8 @@ func TestSubscribeWithArgs(t *testing.T) {
|
||||
require.Equal(t, 1, s.NumClients())
|
||||
require.Equal(t, 1, s.NumClientSubscriptions(clientID))
|
||||
|
||||
require.NoError(t, s.Publish(ctx, "Ka-Zar"))
|
||||
sub.mustReceive(ctx, "Ka-Zar")
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Ka-Zar")))
|
||||
sub.mustReceive(ctx, pubstring("Ka-Zar"))
|
||||
})
|
||||
t.Run("PositiveLimit", func(t *testing.T) {
|
||||
sub := newTestSub(t).must(s.SubscribeWithArgs(ctx, pubsub.SubscribeArgs{
|
||||
@@ -43,8 +50,8 @@ func TestSubscribeWithArgs(t *testing.T) {
|
||||
Query: query.All,
|
||||
Limit: 10,
|
||||
}))
|
||||
require.NoError(t, s.Publish(ctx, "Aggamon"))
|
||||
sub.mustReceive(ctx, "Aggamon")
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Aggamon")))
|
||||
sub.mustReceive(ctx, pubstring("Aggamon"))
|
||||
})
|
||||
}
|
||||
|
||||
@@ -63,7 +70,7 @@ func TestObserver(t *testing.T) {
|
||||
return nil
|
||||
}))
|
||||
|
||||
const input = "Lions and tigers and bears, oh my!"
|
||||
const input = pubstring("Lions and tigers and bears, oh my!")
|
||||
require.NoError(t, s.Publish(ctx, input))
|
||||
<-done
|
||||
require.Equal(t, got, input)
|
||||
@@ -98,14 +105,14 @@ func TestPublishDoesNotBlock(t *testing.T) {
|
||||
go func() {
|
||||
defer close(published)
|
||||
|
||||
require.NoError(t, s.Publish(ctx, "Quicksilver"))
|
||||
require.NoError(t, s.Publish(ctx, "Asylum"))
|
||||
require.NoError(t, s.Publish(ctx, "Ivan"))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Quicksilver")))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Asylum")))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Ivan")))
|
||||
}()
|
||||
|
||||
select {
|
||||
case <-published:
|
||||
sub.mustReceive(ctx, "Quicksilver")
|
||||
sub.mustReceive(ctx, pubstring("Quicksilver"))
|
||||
sub.mustFail(ctx, pubsub.ErrTerminated)
|
||||
case <-time.After(3 * time.Second):
|
||||
t.Fatal("Publishing should not have blocked")
|
||||
@@ -141,13 +148,13 @@ func TestSlowSubscriber(t *testing.T) {
|
||||
Query: query.All,
|
||||
}))
|
||||
|
||||
require.NoError(t, s.Publish(ctx, "Fat Cobra"))
|
||||
require.NoError(t, s.Publish(ctx, "Viper"))
|
||||
require.NoError(t, s.Publish(ctx, "Black Panther"))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Fat Cobra")))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Viper")))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Black Panther")))
|
||||
|
||||
// We had capacity for one item, so we should get that item, but after that
|
||||
// the subscription should have been terminated by the publisher.
|
||||
sub.mustReceive(ctx, "Fat Cobra")
|
||||
sub.mustReceive(ctx, pubstring("Fat Cobra"))
|
||||
sub.mustFail(ctx, pubsub.ErrTerminated)
|
||||
}
|
||||
|
||||
@@ -168,8 +175,8 @@ func TestDifferentClients(t *testing.T) {
|
||||
Attributes: []abci.EventAttribute{{Key: "type", Value: "NewBlock"}},
|
||||
}}
|
||||
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Iceman", events))
|
||||
sub1.mustReceive(ctx, "Iceman")
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Iceman"), events))
|
||||
sub1.mustReceive(ctx, pubstring("Iceman"))
|
||||
|
||||
sub2 := newTestSub(t).must(s.SubscribeWithArgs(ctx, pubsub.SubscribeArgs{
|
||||
ClientID: "client-2",
|
||||
@@ -187,9 +194,9 @@ func TestDifferentClients(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Ultimo", events))
|
||||
sub1.mustReceive(ctx, "Ultimo")
|
||||
sub2.mustReceive(ctx, "Ultimo")
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Ultimo"), events))
|
||||
sub1.mustReceive(ctx, pubstring("Ultimo"))
|
||||
sub2.mustReceive(ctx, pubstring("Ultimo"))
|
||||
|
||||
sub3 := newTestSub(t).must(s.SubscribeWithArgs(ctx, pubsub.SubscribeArgs{
|
||||
ClientID: "client-3",
|
||||
@@ -202,7 +209,7 @@ func TestDifferentClients(t *testing.T) {
|
||||
Attributes: []abci.EventAttribute{{Key: "type", Value: "NewRoundStep"}},
|
||||
}}
|
||||
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Valeria Richards", events))
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Valeria Richards"), events))
|
||||
sub3.mustTimeOut(ctx, 100*time.Millisecond)
|
||||
}
|
||||
|
||||
@@ -215,11 +222,11 @@ func TestSubscribeDuplicateKeys(t *testing.T) {
|
||||
|
||||
testCases := []struct {
|
||||
query string
|
||||
expected interface{}
|
||||
expected types.EventData
|
||||
}{
|
||||
{`withdraw.rewards='17'`, "Iceman"},
|
||||
{`withdraw.rewards='22'`, "Iceman"},
|
||||
{`withdraw.rewards='1' AND withdraw.rewards='22'`, "Iceman"},
|
||||
{`withdraw.rewards='17'`, pubstring("Iceman")},
|
||||
{`withdraw.rewards='22'`, pubstring("Iceman")},
|
||||
{`withdraw.rewards='1' AND withdraw.rewards='22'`, pubstring("Iceman")},
|
||||
{`withdraw.rewards='100'`, nil},
|
||||
}
|
||||
|
||||
@@ -251,7 +258,7 @@ func TestSubscribeDuplicateKeys(t *testing.T) {
|
||||
},
|
||||
}
|
||||
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Iceman", events))
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Iceman"), events))
|
||||
|
||||
if tc.expected != nil {
|
||||
sub.mustReceive(ctx, tc.expected)
|
||||
@@ -280,8 +287,8 @@ func TestClientSubscribesTwice(t *testing.T) {
|
||||
Query: q,
|
||||
}))
|
||||
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Goblin Queen", events))
|
||||
sub1.mustReceive(ctx, "Goblin Queen")
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Goblin Queen"), events))
|
||||
sub1.mustReceive(ctx, pubstring("Goblin Queen"))
|
||||
|
||||
// Subscribing a second time with the same client ID and query fails.
|
||||
{
|
||||
@@ -294,8 +301,8 @@ func TestClientSubscribesTwice(t *testing.T) {
|
||||
}
|
||||
|
||||
// The attempt to re-subscribe does not disrupt the existing sub.
|
||||
require.NoError(t, s.PublishWithEvents(ctx, "Spider-Man", events))
|
||||
sub1.mustReceive(ctx, "Spider-Man")
|
||||
require.NoError(t, s.PublishWithEvents(ctx, pubstring("Spider-Man"), events))
|
||||
sub1.mustReceive(ctx, pubstring("Spider-Man"))
|
||||
}
|
||||
|
||||
func TestUnsubscribe(t *testing.T) {
|
||||
@@ -317,7 +324,7 @@ func TestUnsubscribe(t *testing.T) {
|
||||
}))
|
||||
|
||||
// Publishing should still work.
|
||||
require.NoError(t, s.Publish(ctx, "Nick Fury"))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Nick Fury")))
|
||||
|
||||
// The unsubscribed subscriber should report as such.
|
||||
sub.mustFail(ctx, pubsub.ErrUnsubscribed)
|
||||
@@ -365,8 +372,8 @@ func TestResubscribe(t *testing.T) {
|
||||
|
||||
sub := newTestSub(t).must(s.SubscribeWithArgs(ctx, args))
|
||||
|
||||
require.NoError(t, s.Publish(ctx, "Cable"))
|
||||
sub.mustReceive(ctx, "Cable")
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Cable")))
|
||||
sub.mustReceive(ctx, pubstring("Cable"))
|
||||
}
|
||||
|
||||
func TestUnsubscribeAll(t *testing.T) {
|
||||
@@ -386,7 +393,7 @@ func TestUnsubscribeAll(t *testing.T) {
|
||||
}))
|
||||
|
||||
require.NoError(t, s.UnsubscribeAll(ctx, clientID))
|
||||
require.NoError(t, s.Publish(ctx, "Nick Fury"))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Nick Fury")))
|
||||
|
||||
sub1.mustFail(ctx, pubsub.ErrUnsubscribed)
|
||||
sub2.mustFail(ctx, pubsub.ErrUnsubscribed)
|
||||
@@ -402,13 +409,13 @@ func TestBufferCapacity(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
require.NoError(t, s.Publish(ctx, "Nighthawk"))
|
||||
require.NoError(t, s.Publish(ctx, "Sage"))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Nighthawk")))
|
||||
require.NoError(t, s.Publish(ctx, pubstring("Sage")))
|
||||
|
||||
ctx, cancel = context.WithTimeout(ctx, 100*time.Millisecond)
|
||||
defer cancel()
|
||||
|
||||
require.ErrorIs(t, s.Publish(ctx, "Ironclad"), context.DeadlineExceeded)
|
||||
require.ErrorIs(t, s.Publish(ctx, pubstring("Ironclad")), context.DeadlineExceeded)
|
||||
}
|
||||
|
||||
func newTestServer(ctx context.Context, t testing.TB, logger log.Logger) *pubsub.Server {
|
||||
@@ -436,7 +443,7 @@ func (s *testSub) must(sub *pubsub.Subscription, err error) *testSub {
|
||||
return s
|
||||
}
|
||||
|
||||
func (s *testSub) mustReceive(ctx context.Context, want interface{}) {
|
||||
func (s *testSub) mustReceive(ctx context.Context, want types.EventData) {
|
||||
s.t.Helper()
|
||||
got, err := s.Next(ctx)
|
||||
require.NoError(s.t, err)
|
||||
|
||||
Reference in New Issue
Block a user