2018-08-24 03:17:32 +00:00
|
|
|
package discovery
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
2018-10-08 05:24:39 +00:00
|
|
|
"fmt"
|
2018-08-24 03:17:32 +00:00
|
|
|
"sync"
|
|
|
|
"testing"
|
|
|
|
"time"
|
|
|
|
|
2018-10-08 05:24:39 +00:00
|
|
|
"github.com/ethereum/go-ethereum/event"
|
2018-08-24 03:17:32 +00:00
|
|
|
ma "github.com/multiformats/go-multiaddr"
|
|
|
|
"github.com/status-im/rendezvous"
|
|
|
|
"github.com/stretchr/testify/require"
|
|
|
|
)
|
|
|
|
|
|
|
|
func TestProxyToRendezvous(t *testing.T) {
|
|
|
|
var (
|
|
|
|
topic = "test"
|
|
|
|
id = 101
|
2018-10-08 05:24:39 +00:00
|
|
|
limited = 102
|
|
|
|
limit = 1
|
2018-08-24 03:17:32 +00:00
|
|
|
reg = newRegistry()
|
|
|
|
original = &fake{id: 110, registry: reg, started: true}
|
|
|
|
srv = makeTestRendezvousServer(t, "/ip4/127.0.0.1/tcp/7788")
|
|
|
|
stop = make(chan struct{})
|
2018-10-08 05:24:39 +00:00
|
|
|
feed = &event.Feed{}
|
|
|
|
liveness = 100 * time.Millisecond
|
2018-08-24 03:17:32 +00:00
|
|
|
wg sync.WaitGroup
|
|
|
|
)
|
2018-09-24 05:47:24 +00:00
|
|
|
client, err := rendezvous.NewEphemeral()
|
2018-08-24 03:17:32 +00:00
|
|
|
require.NoError(t, err)
|
|
|
|
reg.Add(topic, id)
|
2018-10-08 05:24:39 +00:00
|
|
|
reg.Add(topic, limited)
|
2018-08-24 03:17:32 +00:00
|
|
|
wg.Add(1)
|
2018-10-08 05:24:39 +00:00
|
|
|
events := make(chan proxyEvent, 10)
|
|
|
|
sub := feed.Subscribe(events)
|
|
|
|
defer sub.Unsubscribe()
|
2018-08-24 03:17:32 +00:00
|
|
|
go func() {
|
|
|
|
defer wg.Done()
|
2018-10-08 05:24:39 +00:00
|
|
|
require.NoError(t, ProxyToRendezvous(original, stop, feed, ProxyOptions{
|
|
|
|
Topic: topic,
|
|
|
|
Servers: []ma.Multiaddr{srv.Addr()},
|
|
|
|
Limit: limit,
|
|
|
|
LivenessWindow: liveness,
|
|
|
|
}))
|
2018-08-24 03:17:32 +00:00
|
|
|
}()
|
2018-10-08 05:24:39 +00:00
|
|
|
require.NoError(t, Consistently(func() (bool, error) {
|
|
|
|
records, err := client.Discover(context.TODO(), srv.Addr(), topic, 10)
|
|
|
|
if err != nil && len(records) < limit {
|
|
|
|
return true, nil
|
|
|
|
}
|
|
|
|
if len(records) > limit {
|
|
|
|
return false, fmt.Errorf("more records than expected: %d != %d", len(records), limit)
|
|
|
|
}
|
|
|
|
var proxied Proxied
|
|
|
|
if err := records[0].Load(&proxied); err != nil {
|
|
|
|
return false, err
|
|
|
|
}
|
|
|
|
if proxied[0] != byte(id) {
|
|
|
|
return false, fmt.Errorf("returned %v instead of %v", proxied[0], id)
|
|
|
|
}
|
|
|
|
return true, nil
|
|
|
|
}, time.Second, 100*time.Millisecond))
|
|
|
|
close(stop)
|
|
|
|
wg.Wait()
|
|
|
|
eventSlice := []proxyEvent{}
|
|
|
|
func() {
|
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case e := <-events:
|
|
|
|
eventSlice = append(eventSlice, e)
|
|
|
|
default:
|
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}()
|
|
|
|
require.Len(t, eventSlice, 2)
|
|
|
|
require.Equal(t, byte(id), eventSlice[0].ID[0])
|
|
|
|
require.Equal(t, proxyStart, eventSlice[0].Type)
|
|
|
|
require.Equal(t, byte(id), eventSlice[1].ID[0])
|
|
|
|
require.Equal(t, proxyStop, eventSlice[1].Type)
|
|
|
|
require.True(t, eventSlice[1].Time.Sub(eventSlice[0].Time) > liveness)
|
|
|
|
}
|
|
|
|
|
|
|
|
func Consistently(f func() (bool, error), timeout, period time.Duration) (err error) {
|
|
|
|
timer := time.After(timeout)
|
|
|
|
ticker := time.Tick(period)
|
|
|
|
var cont bool
|
2018-08-24 03:17:32 +00:00
|
|
|
for {
|
|
|
|
select {
|
|
|
|
case <-timer:
|
2018-10-08 05:24:39 +00:00
|
|
|
return err
|
2018-08-24 03:17:32 +00:00
|
|
|
case <-ticker:
|
2018-10-08 05:24:39 +00:00
|
|
|
cont, err = f()
|
|
|
|
if cont {
|
2018-08-24 03:17:32 +00:00
|
|
|
continue
|
|
|
|
}
|
2018-10-08 05:24:39 +00:00
|
|
|
if err != nil {
|
|
|
|
return err
|
2018-08-24 03:17:32 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|