| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | // Copyright 2016 The go-ethereum Authors | 
					
						
							|  |  |  | // This file is part of the go-ethereum library. | 
					
						
							|  |  |  | // | 
					
						
							|  |  |  | // The go-ethereum library is free software: you can redistribute it and/or modify | 
					
						
							|  |  |  | // it under the terms of the GNU Lesser General Public License as published by | 
					
						
							|  |  |  | // the Free Software Foundation, either version 3 of the License, or | 
					
						
							|  |  |  | // (at your option) any later version. | 
					
						
							|  |  |  | // | 
					
						
							|  |  |  | // The go-ethereum library is distributed in the hope that it will be useful, | 
					
						
							|  |  |  | // but WITHOUT ANY WARRANTY; without even the implied warranty of | 
					
						
							|  |  |  | // MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the | 
					
						
							|  |  |  | // GNU Lesser General Public License for more details. | 
					
						
							|  |  |  | // | 
					
						
							|  |  |  | // You should have received a copy of the GNU Lesser General Public License | 
					
						
							|  |  |  | // along with the go-ethereum library. If not, see <http://www.gnu.org/licenses/>. | 
					
						
							| 
									
										
										
										
											2016-04-14 18:18:24 +02:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | package rpc | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | import ( | 
					
						
							|  |  |  | 	"encoding/json" | 
					
						
							|  |  |  | 	"net" | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	"sync" | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 	"testing" | 
					
						
							|  |  |  | 	"time" | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	"golang.org/x/net/context" | 
					
						
							|  |  |  | ) | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | type NotificationTestService struct { | 
					
						
							|  |  |  | 	mu           sync.Mutex | 
					
						
							|  |  |  | 	unsubscribed bool | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	gotHangSubscriptionReq  chan struct{} | 
					
						
							|  |  |  | 	unblockHangSubscription chan struct{} | 
					
						
							|  |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-08-04 21:18:13 +02:00
										 |  |  | func (s *NotificationTestService) Echo(i int) int { | 
					
						
							|  |  |  | 	return i | 
					
						
							|  |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | func (s *NotificationTestService) wasUnsubCallbackCalled() bool { | 
					
						
							|  |  |  | 	s.mu.Lock() | 
					
						
							|  |  |  | 	defer s.mu.Unlock() | 
					
						
							|  |  |  | 	return s.unsubscribed | 
					
						
							|  |  |  | } | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 
 | 
					
						
							|  |  |  | func (s *NotificationTestService) Unsubscribe(subid string) { | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	s.mu.Lock() | 
					
						
							|  |  |  | 	s.unsubscribed = true | 
					
						
							|  |  |  | 	s.mu.Unlock() | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | func (s *NotificationTestService) SomeSubscription(ctx context.Context, n, val int) (*Subscription, error) { | 
					
						
							| 
									
										
										
										
											2016-04-15 18:05:24 +02:00
										 |  |  | 	notifier, supported := NotifierFromContext(ctx) | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 	if !supported { | 
					
						
							|  |  |  | 		return nil, ErrNotificationsUnsupported | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	// by explicitly creating an subscription we make sure that the subscription id is send back to the client | 
					
						
							|  |  |  | 	// before the first subscription.Notify is called. Otherwise the events might be send before the response | 
					
						
							|  |  |  | 	// for the eth_subscribe method. | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 	subscription := notifier.CreateSubscription() | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 
 | 
					
						
							|  |  |  | 	go func() { | 
					
						
							| 
									
										
										
										
											2017-01-06 19:44:35 +02:00
										 |  |  | 		// test expects n events, if we begin sending event immediately some events | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 		// will probably be dropped since the subscription ID might not be send to | 
					
						
							|  |  |  | 		// the client. | 
					
						
							|  |  |  | 		time.Sleep(5 * time.Second) | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 		for i := 0; i < n; i++ { | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 			if err := notifier.Notify(subscription.ID, val+i); err != nil { | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 				return | 
					
						
							|  |  |  | 			} | 
					
						
							|  |  |  | 		} | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 
 | 
					
						
							|  |  |  | 		select { | 
					
						
							|  |  |  | 		case <-notifier.Closed(): | 
					
						
							|  |  |  | 			s.mu.Lock() | 
					
						
							|  |  |  | 			s.unsubscribed = true | 
					
						
							|  |  |  | 			s.mu.Unlock() | 
					
						
							|  |  |  | 		case <-subscription.Err(): | 
					
						
							|  |  |  | 			s.mu.Lock() | 
					
						
							|  |  |  | 			s.unsubscribed = true | 
					
						
							|  |  |  | 			s.mu.Unlock() | 
					
						
							|  |  |  | 		} | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 	}() | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	return subscription, nil | 
					
						
							|  |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | // HangSubscription blocks on s.unblockHangSubscription before | 
					
						
							|  |  |  | // sending anything. | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | func (s *NotificationTestService) HangSubscription(ctx context.Context, val int) (*Subscription, error) { | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	notifier, supported := NotifierFromContext(ctx) | 
					
						
							|  |  |  | 	if !supported { | 
					
						
							|  |  |  | 		return nil, ErrNotificationsUnsupported | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	s.gotHangSubscriptionReq <- struct{}{} | 
					
						
							|  |  |  | 	<-s.unblockHangSubscription | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 	subscription := notifier.CreateSubscription() | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	go func() { | 
					
						
							| 
									
										
										
										
											2016-07-27 17:47:46 +02:00
										 |  |  | 		notifier.Notify(subscription.ID, val) | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	}() | 
					
						
							|  |  |  | 	return subscription, nil | 
					
						
							|  |  |  | } | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | func TestNotifications(t *testing.T) { | 
					
						
							|  |  |  | 	server := NewServer() | 
					
						
							|  |  |  | 	service := &NotificationTestService{} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	if err := server.RegisterName("eth", service); err != nil { | 
					
						
							|  |  |  | 		t.Fatalf("unable to register test service %v", err) | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	clientConn, serverConn := net.Pipe() | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	go server.ServeCodec(NewJSONCodec(serverConn), OptionMethodInvocation|OptionSubscriptions) | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	out := json.NewEncoder(clientConn) | 
					
						
							|  |  |  | 	in := json.NewDecoder(clientConn) | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	n := 5 | 
					
						
							|  |  |  | 	val := 12345 | 
					
						
							|  |  |  | 	request := map[string]interface{}{ | 
					
						
							|  |  |  | 		"id":      1, | 
					
						
							|  |  |  | 		"method":  "eth_subscribe", | 
					
						
							|  |  |  | 		"version": "2.0", | 
					
						
							|  |  |  | 		"params":  []interface{}{"someSubscription", n, val}, | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	// create subscription | 
					
						
							|  |  |  | 	if err := out.Encode(request); err != nil { | 
					
						
							|  |  |  | 		t.Fatal(err) | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	var subid string | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	response := jsonSuccessResponse{Result: subid} | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 	if err := in.Decode(&response); err != nil { | 
					
						
							|  |  |  | 		t.Fatal(err) | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	var ok bool | 
					
						
							| 
									
										
										
										
											2017-01-09 11:16:06 +01:00
										 |  |  | 	if _, ok = response.Result.(string); !ok { | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 		t.Fatalf("expected subscription id, got %T", response.Result) | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	for i := 0; i < n; i++ { | 
					
						
							|  |  |  | 		var notification jsonNotification | 
					
						
							|  |  |  | 		if err := in.Decode(¬ification); err != nil { | 
					
						
							|  |  |  | 			t.Fatalf("%v", err) | 
					
						
							|  |  |  | 		} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 		if int(notification.Params.Result.(float64)) != val+i { | 
					
						
							|  |  |  | 			t.Fatalf("expected %d, got %d", val+i, notification.Params.Result) | 
					
						
							|  |  |  | 		} | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | 
 | 
					
						
							|  |  |  | 	clientConn.Close() // causes notification unsubscribe callback to be called | 
					
						
							|  |  |  | 	time.Sleep(1 * time.Second) | 
					
						
							|  |  |  | 
 | 
					
						
							| 
									
										
										
										
											2016-07-12 17:47:15 +02:00
										 |  |  | 	if !service.wasUnsubCallbackCalled() { | 
					
						
							| 
									
										
										
										
											2016-03-29 15:07:40 +02:00
										 |  |  | 		t.Error("unsubscribe callback not called after closing connection") | 
					
						
							|  |  |  | 	} | 
					
						
							|  |  |  | } |