-
Notifications
You must be signed in to change notification settings - Fork 0
/
rpc.go
105 lines (91 loc) · 2.24 KB
/
rpc.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
/* ######################################################################
# Author: (zfly1207@126.com)
# Created Time: 2018-08-02 18:54:22
# File Name: rpc.go
# Description:
####################################################################### */
// TODO 还未考虑安全退出
package client
import (
"fmt"
"log"
"net"
"net/rpc"
"github.com/ant-libs-go/ant-mq/client"
)
type rpcClient struct {
opts client.Options
}
var (
rpcc *rpc.Client
listener net.Listener
subscribers = map[string]func(payload []byte) error{}
)
func publish(requester string, server string, topic string, payload []byte) (sessionId string, err error) {
if rpcc == nil {
rpcc, err = rpc.Dial("tcp", server)
if err != nil {
return "", fmt.Errorf("PUB|FAIL|ConnServerExp|%s", err)
}
}
req := client.Request{}
req.Header.Requester = requester
req.Header.Topic = topic
req.Body = payload
resp := &client.Response{}
err = rpcc.Call("AntMqServer.Pub", req, &resp)
if err != nil {
rpcc = nil
return "", fmt.Errorf("PUB|FAIL|CallApiExp|%s", err)
}
if resp.Header.Status != client.RespStatusSucc {
return "", fmt.Errorf("PUB|FAIL|CallApiExpStatus-%d", resp.Header.Status)
}
sessionId = resp.Header.SessionId
return
}
func (this *rpcClient) Publish(topic string, payload []byte) (err error) {
var sessionId string
for i := int32(1); this.opts.Retry >= i; i++ {
sessionId, err = publish(this.opts.Requester, this.opts.Server, topic, payload)
if err == nil {
log.Printf("PUB|SUCC-%s", sessionId)
break
}
log.Printf("%s|Try-%d", err, i)
}
return
}
func (this *rpcClient) Subscribe(subs map[string]func(payload []byte) error) (err error) {
for k, v := range subs {
subscribers[k] = v
}
rpc.RegisterName("AntMqClient", NewHandler())
listener, err = net.Listen("tcp", this.opts.Listen)
if err != nil {
return fmt.Errorf("SUB|FAIL|Listen-%s|%s", this.opts.Listen, err)
}
for {
conn, err := listener.Accept()
if err != nil {
continue
}
go rpc.ServeConn(conn)
}
return
}
func (this *rpcClient) Close() error {
return nil
}
func New(name string, opts ...client.Option) *rpcClient {
options := client.Options{
Requester: name,
}
for _, opt := range opts {
opt(&options)
}
o := &rpcClient{
opts: options,
}
return o
}