From d50f242b6d9c2c4da40f59480f9c821e26ec9e77 Mon Sep 17 00:00:00 2001 From: nkbai Date: Fri, 22 Sep 2017 09:27:31 +0800 Subject: [PATCH 1/2] rpc: client.go remove duplicate code --- rpc/client.go | 110 ++++++++++++++++++++++---------------------------- 1 file changed, 48 insertions(+), 62 deletions(-) diff --git a/rpc/client.go b/rpc/client.go index f02366a39c..21ec310f72 100644 --- a/rpc/client.go +++ b/rpc/client.go @@ -349,6 +349,52 @@ func (c *Client) BatchCallContext(ctx context.Context, b []BatchElem) error { return err } +// Subscribe calls the "[service]_subscribe" method with the given arguments, +// registering a subscription. Server notifications for the subscription are +// sent to the given channel. The element type of the channel must match the +// expected type of content returned by the subscription. +// +// The context argument cancels the RPC request that sets up the subscription but has no +// effect on the subscription after subscribe has returned. +// +// Slow subscribers will be dropped eventually. Client buffers up to 8000 notifications +// before considering the subscriber dead. The subscription Err channel will receive +// ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure +// that the channel usually has at least one reader to prevent this issue. +func (c *Client) Subscribe(ctx context.Context, service string, channel interface{}, args ...interface{}) (*ClientSubscription, error) { + // Check type of channel first. + chanVal := reflect.ValueOf(channel) + if chanVal.Kind() != reflect.Chan || chanVal.Type().ChanDir()&reflect.SendDir == 0 { + panic("first argument to Subscribe must be a writable channel") + } + if chanVal.IsNil() { + panic("channel given to Subscribe must not be nil") + } + if c.isHTTP { + return nil, ErrNotificationsUnsupported + } + + msg, err := c.newMessage(service+subscribeMethodSuffix, args...) + if err != nil { + return nil, err + } + op := &requestOp{ + ids: []json.RawMessage{msg.ID}, + resp: make(chan *jsonrpcMessage), + sub: newClientSubscription(c, service, chanVal), + } + + // Send the subscription request. + // The arrival and validity of the response is signaled on sub.quit. + if err := c.send(ctx, op, msg); err != nil { + return nil, err + } + if _, err := op.wait(ctx); err != nil { + return nil, err + } + return op.sub, nil +} + // ShhSubscribe calls the "shh_subscribe" method with the given arguments, // registering a subscription. Server notifications for the subscription are // sent to the given channel. The element type of the channel must match the @@ -362,37 +408,7 @@ func (c *Client) BatchCallContext(ctx context.Context, b []BatchElem) error { // ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure // that the channel usually has at least one reader to prevent this issue. func (c *Client) ShhSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) { - // Check type of channel first. - chanVal := reflect.ValueOf(channel) - if chanVal.Kind() != reflect.Chan || chanVal.Type().ChanDir()&reflect.SendDir == 0 { - panic("first argument to ShhSubscribe must be a writable channel") - } - if chanVal.IsNil() { - panic("channel given to ShhSubscribe must not be nil") - } - if c.isHTTP { - return nil, ErrNotificationsUnsupported - } - - msg, err := c.newMessage("shh"+subscribeMethodSuffix, args...) - if err != nil { - return nil, err - } - op := &requestOp{ - ids: []json.RawMessage{msg.ID}, - resp: make(chan *jsonrpcMessage), - sub: newClientSubscription(c, "shh", chanVal), - } - - // Send the subscription request. - // The arrival and validity of the response is signaled on sub.quit. - if err := c.send(ctx, op, msg); err != nil { - return nil, err - } - if _, err := op.wait(ctx); err != nil { - return nil, err - } - return op.sub, nil + return c.Subscribe(ctx, "shh", channel, args...) } // EthSubscribe calls the "eth_subscribe" method with the given arguments, @@ -408,37 +424,7 @@ func (c *Client) ShhSubscribe(ctx context.Context, channel interface{}, args ... // ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure // that the channel usually has at least one reader to prevent this issue. func (c *Client) EthSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) { - // Check type of channel first. - chanVal := reflect.ValueOf(channel) - if chanVal.Kind() != reflect.Chan || chanVal.Type().ChanDir()&reflect.SendDir == 0 { - panic("first argument to EthSubscribe must be a writable channel") - } - if chanVal.IsNil() { - panic("channel given to EthSubscribe must not be nil") - } - if c.isHTTP { - return nil, ErrNotificationsUnsupported - } - - msg, err := c.newMessage("eth"+subscribeMethodSuffix, args...) - if err != nil { - return nil, err - } - op := &requestOp{ - ids: []json.RawMessage{msg.ID}, - resp: make(chan *jsonrpcMessage), - sub: newClientSubscription(c, "eth", chanVal), - } - - // Send the subscription request. - // The arrival and validity of the response is signaled on sub.quit. - if err := c.send(ctx, op, msg); err != nil { - return nil, err - } - if _, err := op.wait(ctx); err != nil { - return nil, err - } - return op.sub, nil + return c.Subscribe(ctx, "eth", channel, args...) } func (c *Client) newMessage(method string, paramsIn ...interface{}) (*jsonrpcMessage, error) { From d676110959c3ef3dc2ec000a5f365d700bd3b7ef Mon Sep 17 00:00:00 2001 From: nkbai Date: Fri, 22 Sep 2017 11:06:56 +0800 Subject: [PATCH 2/2] whisper: add client test case * whisper: fix api Version error --- whisper/shhclient/client.go | 4 +- whisper/shhclient/client_test.go | 239 +++++++++++++++++++++++++++++++ 2 files changed, 241 insertions(+), 2 deletions(-) create mode 100644 whisper/shhclient/client_test.go diff --git a/whisper/shhclient/client.go b/whisper/shhclient/client.go index 61c1b7ab72..61b4775d95 100644 --- a/whisper/shhclient/client.go +++ b/whisper/shhclient/client.go @@ -45,8 +45,8 @@ func NewClient(c *rpc.Client) *Client { } // Version returns the Whisper sub-protocol version. -func (sc *Client) Version(ctx context.Context) (uint, error) { - var result uint +func (sc *Client) Version(ctx context.Context) (string, error) { + var result string err := sc.c.CallContext(ctx, &result, "shh_version") return result, err } diff --git a/whisper/shhclient/client_test.go b/whisper/shhclient/client_test.go new file mode 100644 index 0000000000..f5b7214536 --- /dev/null +++ b/whisper/shhclient/client_test.go @@ -0,0 +1,239 @@ +package shhclient + +import ( + "bytes" + "context" + "github.com/ethereum/go-ethereum/node" + "github.com/ethereum/go-ethereum/whisper/whisperv5" + "testing" + "time" +) + +func getIPCPath() string { + return node.DefaultIPCEndpoint("geth") +} +func getClient() (*Client, error) { + return Dial(getIPCPath()) +} + +func TestBasic(t *testing.T) { + var id string = "test" + c, err := getClient() + if err != nil { //if geth not start,just skip test. + t.Log("skip Client test because of that geth is not started") + return + } + defer func() { + c.c.Close() + }() + ctx := context.Background() + if err != nil { + t.Error(err) + return + } + version, err := c.Version(ctx) + if err != nil { + t.Error(err) + return + } + if version != whisperv5.ProtocolVersionStr { + t.Fatalf("wrong version: %d.", version) + } + _, err = c.Info(ctx) + if err != nil { + t.Error(err) + return + } + exist, err := c.HasKeyPair(ctx, id) + if err != nil { + t.Fatalf("failed initial HasIdentity: %s.", err) + } + if exist { + t.Fatalf("failed initial HasIdentity: false positive.") + } + id = "arbitrary text" + id2 := "another arbitrary string" + + exist, err = c.HasSymmetricKey(ctx, id) + if err != nil { + t.Fatalf("failed HasSymKey: %s.", err) + } + if exist { + t.Fatalf("failed HasSymKey: false positive.") + } + + id, err = c.NewSymmetricKey(ctx) + if err != nil { + t.Fatalf("failed GenerateSymKey: %s.", err) + } + + exist, err = c.HasSymmetricKey(ctx, id) + if err != nil { + t.Fatalf("failed HasSymKey(): %s.", err) + } + if !exist { + t.Fatalf("failed HasSymKey(): false negative.") + } + + var password = []byte("some stuff here") + id, err = c.GenerateSymmetricKeyFromPassword(ctx, password) + if err != nil { + t.Fatalf("failed AddSymKey: %s.", err) + } + + id2, err = c.GenerateSymmetricKeyFromPassword(ctx, password) + if err != nil { + t.Fatalf("failed AddSymKey: %s.", err) + } + + exist, err = c.HasSymmetricKey(ctx, id2) + if err != nil { + t.Fatalf("failed HasSymKey(id2): %s.", err) + } + if !exist { + t.Fatalf("failed HasSymKey(id2): false negative.") + } + + k1, err := c.GetSymmetricKey(ctx, id) + if err != nil { + t.Fatalf("failed GetSymKey(id): %s.", err) + } + k2, err := c.GetSymmetricKey(ctx, id2) + if err != nil { + t.Fatalf("failed GetSymKey(id2): %s.", err) + } + + if !bytes.Equal(k1, k2) { + t.Fatalf("installed keys are not equal") + } + + err = c.DeleteSymmetricKey(ctx, id) + if err != nil { + t.Fatalf("failed DeleteSymKey(id): %s.", err) + } + + exist, err = c.HasSymmetricKey(ctx, id) + if err != nil { + t.Fatalf("failed HasSymKey(id): %s.", err) + } + if exist { + t.Fatalf("failed HasSymKey(id): false positive.") + } +} + +func TestSubscribe(t *testing.T) { + var err error + var ctx = context.Background() + + c, err := getClient() + if err != nil { //if geth not start,just skip test. + t.Log("skip Client test because of that geth is not started") + return + } + defer func() { + c.c.Close() + }() + symKeyID, err := c.NewSymmetricKey(ctx) + if err != nil { + t.Fatalf("failed to GenerateSymKey: %s.", err) + } + + var f whisperv5.Criteria + f.SymKeyID = symKeyID + f.Topics = make([]whisperv5.TopicType, 2) + f.Topics[0] = whisperv5.TopicType{0xf8, 0xe9, 0xa0, 0xba} + f.Topics[1] = whisperv5.TopicType{0xcb, 0x3c, 0xdd, 0xee} + ch := make(chan *whisperv5.Message) + sub, err := c.SubscribeMessages(ctx, f, ch) + if err != nil { + t.Fatalf("failed to subscribe: %s.", err) + } + sub.Unsubscribe() + close(ch) +} + +func TestIntegrationSymWithFilter(t *testing.T) { + c, err := getClient() + if err != nil { + t.Log("skip Client test because of that geth is not started") + return + } + defer func() { + c.c.Close() + }() + ctx := context.Background() + symKeyID, err := c.NewSymmetricKey(ctx) + if err != nil { + t.Fatalf("failed to GenerateSymKey: %s.", err) + } + + sigKeyID, err := c.NewKeyPair(ctx) + if err != nil { + t.Fatalf("failed NewIdentity: %s.", err) + } + if len(sigKeyID) == 0 { + t.Fatalf("wrong signature.") + } + + exist, err := c.HasKeyPair(ctx, sigKeyID) + if err != nil { + t.Fatalf("failed HasIdentity: %s.", err) + } + if !exist { + t.Fatalf("failed HasIdentity: does not exist.") + } + + sigPubKey, err := c.PublicKey(ctx, sigKeyID) + if err != nil { + t.Fatalf("failed GetPublicKey: %s.", err) + } + + var topics [2]whisperv5.TopicType + topics[0] = whisperv5.TopicType{0x00, 0x7f, 0x80, 0xff} + topics[1] = whisperv5.TopicType{0xf2, 0x6e, 0x77, 0x79} + var f whisperv5.Criteria + f.SymKeyID = symKeyID + f.Topics = make([]whisperv5.TopicType, 2) + f.Topics[0] = topics[0] + f.Topics[1] = topics[1] + f.MinPow = whisperv5.DefaultMinimumPoW / 2 + f.Sig = sigPubKey + f.AllowP2P = false + ch := make(chan *whisperv5.Message) + sub, err := c.SubscribeMessages(ctx, f, ch) + if err != nil { + t.Fatalf("failed to create new filter: %s.", err) + } + defer func() { + sub.Unsubscribe() + close(ch) + }() + var p whisperv5.NewMessage + p.TTL = 1 + p.SymKeyID = symKeyID + p.Sig = sigKeyID + p.Padding = []byte("test string") + p.Payload = []byte("extended test string") + p.PowTarget = whisperv5.DefaultMinimumPoW + p.PowTime = 2 + p.Topic = whisperv5.TopicType{0xf2, 0x6e, 0x77, 0x79} + + err = c.Post(ctx, p) + if err != nil { + t.Fatalf("failed to post message: %s.", err) + } + + var msg *whisperv5.Message = nil + select { + case <-time.After(time.Second * 10): + case msg = <-ch: + } + if msg == nil { + t.Fatalf("failed to GetFilterChanges: got messages.") + } + + text := string(msg.Payload) + if text != string("extended test string") { + t.Fatalf("failed to decrypt first message: %s.", text) + } +}