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 balancer/group_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ func TestGroup_Get(t *testing.T) {
called++
return balancer.NewRoundRobinBalancer()
})
defer g.Close()
for range [2]struct{}{} {
b := g.Get(fakeGroupId)
assert.NotNil(t, b)
Expand Down
14 changes: 9 additions & 5 deletions core/transport/tcp_transport.go
Original file line number Diff line number Diff line change
Expand Up @@ -151,16 +151,20 @@ func NewTCPClientTransport(c net.Conn) *Transport {
}

// NewTCPClientTransportWithAddr creates a new transport.
func NewTCPClientTransportWithAddr(network, addr string, tlsConfig *tls.Config) (tp *Transport, err error) {
var rawConn net.Conn
func NewTCPClientTransportWithAddr(ctx context.Context, network, addr string, tlsConfig *tls.Config) (tp *Transport, err error) {
var conn net.Conn
if tlsConfig == nil {
rawConn, err = net.Dial(network, addr)
var dial net.Dialer
conn, err = dial.DialContext(ctx, network, addr)
} else {
rawConn, err = tls.Dial(network, addr, tlsConfig)
dial := tls.Dialer{
Config: tlsConfig,
}
conn, err = dial.DialContext(ctx, network, addr)
}
if err != nil {
return
}
tp = NewTCPClientTransport(rawConn)
tp = NewTCPClientTransport(conn)
return
}
12 changes: 6 additions & 6 deletions core/transport/websocket_transport.go
Original file line number Diff line number Diff line change
Expand Up @@ -177,20 +177,20 @@ func NewWebsocketServerTransportWithAddr(addr string, path string, upgrader *web
}

// NewWebsocketClientTransport creates a new client-side transport.
func NewWebsocketClientTransport(url string, config *tls.Config, header http.Header) (*Transport, error) {
var d *websocket.Dialer
func NewWebsocketClientTransport(ctx context.Context, url string, config *tls.Config, header http.Header) (*Transport, error) {
var dial *websocket.Dialer
if config == nil {
d = websocket.DefaultDialer
dial = websocket.DefaultDialer
} else {
d = &websocket.Dialer{
dial = &websocket.Dialer{
Proxy: http.ProxyFromEnvironment,
HandshakeTimeout: 45 * time.Second,
TLSClientConfig: config,
}
}
wsConn, _, err := d.Dial(url, header)
conn, _, err := dial.DialContext(ctx, url, header)
if err != nil {
return nil, errors.Wrap(err, "dial websocket failed")
}
return NewTransport(NewWebsocketConnection(wsConn)), nil
return NewTransport(NewWebsocketConnection(conn)), nil
}
Binary file removed examples/echo/echo
Binary file not shown.
5 changes: 1 addition & 4 deletions examples/echo_bench/echo_bench.go
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,6 @@ import (
"sync"
"time"

"github.com/jjeffcaii/reactor-go/scheduler"
"github.com/rsocket/rsocket-go"
"github.com/rsocket/rsocket-go/core/transport"
"github.com/rsocket/rsocket-go/payload"
Expand Down Expand Up @@ -47,7 +46,6 @@ func main() {
rand.Read(data)

now := time.Now()
ctx := context.Background()

sub := rx.NewSubscriber(
rx.OnNext(func(input payload.Payload) error {
Expand All @@ -57,9 +55,8 @@ func main() {
return nil
}),
)

for i := 0; i < n; i++ {
client.RequestResponse(payload.New(data, nil)).SubscribeOn(scheduler.Parallel()).SubscribeWith(ctx, sub)
client.RequestResponse(payload.New(data, nil)).SubscribeWith(context.Background(), sub)
}
wg.Wait()
cost := time.Since(now)
Expand Down
2 changes: 1 addition & 1 deletion go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ require (
github.com/golang/mock v1.4.3
github.com/google/uuid v1.1.1
github.com/gorilla/websocket v1.4.1
github.com/jjeffcaii/reactor-go v0.2.3
github.com/jjeffcaii/reactor-go v0.2.4
github.com/pkg/errors v0.9.1
github.com/stretchr/testify v1.4.0
github.com/urfave/cli/v2 v2.1.1
Expand Down
7 changes: 2 additions & 5 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@ github.com/google/uuid v1.1.1 h1:Gkbcsh/GbpXz7lPftLA3P6TYMwjCLYm83jiFQZF/3gY=
github.com/google/uuid v1.1.1/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo=
github.com/gorilla/websocket v1.4.1 h1:q7AeDBpnBk8AogcD4DSag/Ukw/KV+YhzLj2bP5HvKCM=
github.com/gorilla/websocket v1.4.1/go.mod h1:YR8l580nyteQvAITg2hZ9XVh4b55+EU/adAjf1fMHhE=
github.com/jjeffcaii/reactor-go v0.2.3 h1:qkcnnNyJ241hq15EMjGXP/mtqxRQlITy7eed2qzZdYQ=
github.com/jjeffcaii/reactor-go v0.2.3/go.mod h1:I4qZrpZcsqjzo3pjq0XWGBTpdFXB95XeYinrPYETNL4=
github.com/jjeffcaii/reactor-go v0.2.4 h1:Q3N/0Ngt1Ywi7ezye2LQ+mU1vNdHxyG5ZRk3W2EWmYA=
github.com/jjeffcaii/reactor-go v0.2.4/go.mod h1:I4qZrpZcsqjzo3pjq0XWGBTpdFXB95XeYinrPYETNL4=
github.com/panjf2000/ants/v2 v2.4.1 h1:7RtUqj5lGOw0WnZhSKDZ2zzJhaX5490ZW1sUolRXCxY=
github.com/panjf2000/ants/v2 v2.4.1/go.mod h1:f6F0NZVFsGCp5A7QW/Zj/m92atWwOkY0OIhFxRNFr4A=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
Expand All @@ -31,7 +31,6 @@ github.com/urfave/cli/v2 v2.1.1/go.mod h1:SE9GqnLQmjVa0iPEY0f1w3ygNIYcIJ0OKPMoW2
go.uber.org/atomic v1.5.1 h1:rsqfU5vBkVknbhUGbAUwQKR2H4ItV8tjJ+6kJX4cxHM=
go.uber.org/atomic v1.5.1/go.mod h1:sABNBOSYdrvTF6hTgEIbc7YasKWGhgEQZyfxyTvoXHQ=
golang.org/x/crypto v0.0.0-20190308221718-c2843e01d9a2/go.mod h1:djNgcEr1/C05ACkg1iLfiJU5Ep61QUkGW8qpdssI0+w=
golang.org/x/lint v0.0.0-20190930215403-16217165b5de h1:5hukYrvBGR8/eNkX5mdUezrA6JiaEZDtJb9Ei+1LlBs=
golang.org/x/lint v0.0.0-20190930215403-16217165b5de/go.mod h1:6SW0HCj/g11FgYtHlgUYUwCkIfeOF89ocIRzGO/8vkc=
golang.org/x/net v0.0.0-20190311183353-d8887717615a/go.mod h1:t9HGtf8HONx5eT2rtn7q6eTqICYqUVnKs3thJo3Qplg=
golang.org/x/net v0.0.0-20190620200207-3b0461eec859/go.mod h1:z5CRVTTTmAJ677TzLLGU+0bjPO0LkuOLi4/5GtJWs/s=
Expand All @@ -41,10 +40,8 @@ golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fq
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/tools v0.0.0-20190311212946-11955173bddd/go.mod h1:LCzVGOaR6xXOjkQ3onu1FJEFr0SW1gC7cKk1uF8kGRs=
golang.org/x/tools v0.0.0-20190425150028-36563e24a262/go.mod h1:RgjU9mgBXZiqYHBnxXauZ1Gv1EHHAz9KjViQ78xBX0Q=
golang.org/x/tools v0.0.0-20191029041327-9cc4af7d6b2c h1:IGkKhmfzcztjm6gYkykvu/NiS8kaqbCWAEWWAyf8J5U=
golang.org/x/tools v0.0.0-20191029041327-9cc4af7d6b2c/go.mod h1:b+2E5dAYhXwXZwtnZ6UAqBI28+e2cm9otk0dWdXHAEo=
golang.org/x/xerrors v0.0.0-20190717185122-a985d3407aa7/go.mod h1:I/5z698sn9Ka8TeJc9MKroUUfqBBauWjQqLJ2OPfmY0=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405 h1:yhCVgyC4o1eVCa2tZl7eS0r+SDo693bJlVdllGtEeKM=
gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0=
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
gopkg.in/yaml.v2 v2.2.7 h1:VUgggvou5XRW9mHwD/yXxIYSMtY0zoKQf/v226p2nyo=
Expand Down
8 changes: 4 additions & 4 deletions internal/socket/duplex.go
Original file line number Diff line number Diff line change
Expand Up @@ -355,7 +355,7 @@ func (dc *DuplexConnection) RequestChannel(publisher rx.Publisher) (ret flux.Flu
})
return
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
dc.register(sid, requestChannelCallback{rcv: receiving, snd: s})
s.Request(1)
}),
Expand Down Expand Up @@ -419,7 +419,7 @@ func (dc *DuplexConnection) respondRequestResponse(receiving fragmentation.Heade
rx.OnError(func(e error) {
dc.writeError(sid, e)
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
dc.register(sid, requestResponseCallbackReverse{su: s})
s.Request(rx.RequestMax)
}),
Expand Down Expand Up @@ -505,7 +505,7 @@ func (dc *DuplexConnection) respondRequestChannel(pl fragmentation.HeaderAndPayl
<-complete.DoneNotify()
}
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
dc.register(sid, requestChannelCallbackReverse{rcv: receivingProcessor, snd: s})
close(mustSub)
s.Request(initRequestN)
Expand Down Expand Up @@ -602,7 +602,7 @@ func (dc *DuplexConnection) respondRequestStream(receiving fragmentation.HeaderA
dc.sendPayload(sid, elem, core.FlagNext)
return nil
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
dc.register(sid, requestStreamCallbackReverse{su: s})
s.Request(n32)
}),
Expand Down
4 changes: 2 additions & 2 deletions internal/socket/resumable_client_socket_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,7 +76,7 @@ func TestNewResumableClientSocket(t *testing.T) {
}

result, err := rcs.RequestResponse(payload.New(fakeData, fakeMetadata)).
DoOnSubscribe(func(s rx.Subscription) {
DoOnSubscribe(func(ctx context.Context, s rx.Subscription) {
readChan <- framing.NewPayloadFrame(nextRequestId(), fakeData, fakeMetadata, core.FlagComplete)
}).
Block(context.Background())
Expand All @@ -90,7 +90,7 @@ func TestNewResumableClientSocket(t *testing.T) {
stream = append(stream, input)
return nil
}).
DoOnSubscribe(func(s rx.Subscription) {
DoOnSubscribe(func(ctx context.Context, s rx.Subscription) {
nextId := nextRequestId()
readChan <- framing.NewPayloadFrame(nextId, fakeData, fakeMetadata, core.FlagNext)
readChan <- framing.NewPayloadFrame(nextId, fakeData, fakeMetadata, core.FlagNext)
Expand Down
4 changes: 2 additions & 2 deletions internal/socket/simple_client_socket_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -68,7 +68,7 @@ func TestNewClient(t *testing.T) {
}

result, err := cli.RequestResponse(payload.New(fakeData, fakeMetadata)).
DoOnSubscribe(func(s rx.Subscription) {
DoOnSubscribe(func(ctx context.Context, s rx.Subscription) {
readChan <- framing.NewPayloadFrame(nextRequestId(), fakeData, fakeMetadata, core.FlagComplete)
}).
Block(context.Background())
Expand All @@ -82,7 +82,7 @@ func TestNewClient(t *testing.T) {
stream = append(stream, input)
return nil
}).
DoOnSubscribe(func(s rx.Subscription) {
DoOnSubscribe(func(ctx context.Context, s rx.Subscription) {
nextId := nextRequestId()
readChan <- framing.NewPayloadFrame(nextId, fakeData, fakeMetadata, core.FlagNext)
readChan <- framing.NewPayloadFrame(nextId, fakeData, fakeMetadata, core.FlagNext)
Expand Down
2 changes: 1 addition & 1 deletion rsocket_example_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -133,7 +133,7 @@ func ExampleConnect() {
s.Request(1)
return nil
}).
Subscribe(context.Background(), rx.OnSubscribe(func(s rx.Subscription) {
Subscribe(context.Background(), rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
s.Request(1)
}))
// Simple RequestChannel.
Expand Down
107 changes: 105 additions & 2 deletions rsocket_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,13 @@ func TestResume(t *testing.T) {
}()

go func(ctx context.Context) {
defer func() {
select {
case <-started:
default:
close(started)
}
}()
_ = Receive().
OnStart(func() {
close(started)
Expand Down Expand Up @@ -158,6 +165,13 @@ func TestConnectBroken(t *testing.T) {
port := 8787

go func(ctx context.Context) {
defer func() {
select {
case <-started:
default:
close(started)
}
}()
_ = Receive().
OnStart(func() {
close(started)
Expand Down Expand Up @@ -206,6 +220,15 @@ func TestBiDirection(t *testing.T) {
defer cancel()

go func(ctx context.Context) {

defer func() {
select {
case <-started:
default:
close(started)
}
}()

l, _ := lease.NewSimpleFactory(3*time.Second, 1*time.Second, 1*time.Second, 10)
_ = Receive().
Lease(l).
Expand Down Expand Up @@ -313,6 +336,15 @@ func testAll(t *testing.T, proto string, clientTp transport.ClientTransporter, s
serving := make(chan struct{})

go func(ctx context.Context) {

defer func() {
select {
case <-serving:
default:
close(serving)
}
}()

err := Receive().
Fragment(128).
OnStart(func() {
Expand Down Expand Up @@ -453,7 +485,7 @@ func testRequestStreamOneByOne(ctx context.Context, cli Client, t *testing.T) {
su.Request(1)
return nil
}).
Subscribe(ctx, rx.OnSubscribe(func(s rx.Subscription) {
Subscribe(ctx, rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
su = s
su.Request(1)
}))
Expand Down Expand Up @@ -515,7 +547,7 @@ func testRequestChannelOneByOne(ctx context.Context, cli Client, t *testing.T) {
Subscribe(ctx, rx.OnNext(func(elem payload.Payload) error {
su.Request(1)
return nil
}), rx.OnSubscribe(func(s rx.Subscription) {
}), rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
su = s
su.Request(1)
}))
Expand Down Expand Up @@ -567,3 +599,74 @@ func startProxy(addr string, ch chan net.Listener, upstreamAddr string) {
}

}

type delayedRSocket struct {
}

func (d delayedRSocket) FireAndForget(message payload.Payload) {
panic("implement me")
}

func (d delayedRSocket) MetadataPush(message payload.Payload) {
panic("implement me")
}

func (d delayedRSocket) RequestResponse(message payload.Payload) mono.Mono {
return mono.Create(func(ctx context.Context, sink mono.Sink) {
time.AfterFunc(300*time.Millisecond, func() {
sink.Success(message)
})
})
}

func (d delayedRSocket) RequestStream(message payload.Payload) flux.Flux {
panic("implement me")
}

func (d delayedRSocket) RequestChannel(messages rx.Publisher) flux.Flux {
panic("implement me")
}

func TestContextTimeout(t *testing.T) {
var responder delayedRSocket
started := make(chan struct{})
go func() {
defer func() {
select {
case <-started:
default:
close(started)
}
}()

_ = Receive().
OnStart(func() {
close(started)
}).
Acceptor(func(setup payload.SetupPayload, sendingSocket CloseableRSocket) (RSocket, error) {
return responder, nil
}).
Transport(TCPServer().SetAddr(":8088").Build()).
Serve(context.Background())
}()

<-started

tp := TCPClient().SetAddr("127.0.0.1:8088").Build()

// simulate timeout
ctxMustTimeout, cancel := context.WithTimeout(context.Background(), 1*time.Nanosecond)
defer cancel()
_, err := Connect().Transport(tp).Start(ctxMustTimeout)
assert.Error(t, err, "should connect timeout")

cli, err := Connect().Transport(tp).Start(context.Background())
assert.NoError(t, err, "should connect success")
defer cli.Close()

ctx, cancel2 := context.WithTimeout(context.Background(), 100*time.Millisecond)
defer cancel2()

_, err = cli.RequestResponse(fakeRequest).Block(ctx)
assert.Error(t, err, "should return error")
}
2 changes: 1 addition & 1 deletion rx/flux/flux.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,7 @@ type Flux interface {
// DoOnSubscribe add behavior triggered when the Flux is done being subscribed.
DoOnSubscribe(rx.FnOnSubscribe) Flux
// Map transform the items emitted by this Flux by applying a synchronous function to each item.
Map(func(payload.Payload) (payload.Payload, error)) Flux
Map(rx.FnTransform) Flux
// SwitchOnFirst transform the current Flux once it emits its first element, making a conditional transformation possible.
SwitchOnFirst(FnSwitchOnFirst) Flux
// SubscribeOn run subscribe, onSubscribe and request on a specified scheduler.
Expand Down
6 changes: 3 additions & 3 deletions rx/flux/flux_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -149,7 +149,7 @@ func TestCreate(t *testing.T) {
DoOnComplete(func() {
fmt.Println("complete")
}).
Subscribe(context.Background(), rx.OnSubscribe(func(s rx.Subscription) {
Subscribe(context.Background(), rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
su = s
su.Request(1)
}))
Expand Down Expand Up @@ -234,7 +234,7 @@ func TestFluxRequest(t *testing.T) {
rx.OnComplete(func() {
fmt.Println("complete")
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
su = s
su.Request(1)
fmt.Println("request:", 1)
Expand Down Expand Up @@ -271,7 +271,7 @@ func TestFluxProcessorWithRequest(t *testing.T) {
su.Request(1)
return nil
}),
rx.OnSubscribe(func(s rx.Subscription) {
rx.OnSubscribe(func(ctx context.Context, s rx.Subscription) {
su = s
su.Request(1)
}),
Expand Down
Loading