mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-17 01:13:45 +00:00
rpc: client.go remove duplicate code
This commit is contained in:
parent
d558a595ad
commit
d50f242b6d
1 changed files with 48 additions and 62 deletions
110
rpc/client.go
110
rpc/client.go
|
|
@ -349,6 +349,52 @@ func (c *Client) BatchCallContext(ctx context.Context, b []BatchElem) error {
|
||||||
return err
|
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,
|
// ShhSubscribe calls the "shh_subscribe" method with the given arguments,
|
||||||
// registering a subscription. Server notifications for the subscription are
|
// registering a subscription. Server notifications for the subscription are
|
||||||
// sent to the given channel. The element type of the channel must match the
|
// 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
|
// ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure
|
||||||
// that the channel usually has at least one reader to prevent this issue.
|
// 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) {
|
func (c *Client) ShhSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) {
|
||||||
// Check type of channel first.
|
return c.Subscribe(ctx, "shh", channel, args...)
|
||||||
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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// EthSubscribe calls the "eth_subscribe" method with the given arguments,
|
// 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
|
// ErrSubscriptionQueueOverflow. Use a sufficiently large buffer on the channel or ensure
|
||||||
// that the channel usually has at least one reader to prevent this issue.
|
// 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) {
|
func (c *Client) EthSubscribe(ctx context.Context, channel interface{}, args ...interface{}) (*ClientSubscription, error) {
|
||||||
// Check type of channel first.
|
return c.Subscribe(ctx, "eth", channel, args...)
|
||||||
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
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func (c *Client) newMessage(method string, paramsIn ...interface{}) (*jsonrpcMessage, error) {
|
func (c *Client) newMessage(method string, paramsIn ...interface{}) (*jsonrpcMessage, error) {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue