|
|
|
|
@ -4,6 +4,7 @@ import (
|
|
|
|
|
"context" |
|
|
|
|
"encoding/base64" |
|
|
|
|
"encoding/json" |
|
|
|
|
"errors" |
|
|
|
|
"flag" |
|
|
|
|
"fmt" |
|
|
|
|
"github.com/bytedance/sonic" |
|
|
|
|
@ -16,6 +17,7 @@ import (
|
|
|
|
|
"math/rand" |
|
|
|
|
"net" |
|
|
|
|
"net/http" |
|
|
|
|
"reflect" |
|
|
|
|
"runtime" |
|
|
|
|
"sonet/api/gen/auth" |
|
|
|
|
"sonet/api/gen/chat" |
|
|
|
|
@ -23,6 +25,7 @@ import (
|
|
|
|
|
"sonet/pkg/grpc/discovery" |
|
|
|
|
"sonet/pkg/protocol" |
|
|
|
|
"sonet/pkg/protocol/deliver" |
|
|
|
|
"sonet/pkg/utils/collect" |
|
|
|
|
"sonet/pkg/utils/logger" |
|
|
|
|
"sonet/pkg/utils/security" |
|
|
|
|
"sonet/pkg/utils/shutdown" |
|
|
|
|
@ -34,28 +37,29 @@ import (
|
|
|
|
|
) |
|
|
|
|
|
|
|
|
|
var ( |
|
|
|
|
gatewayHttp = "http://192.168.110.41:7000" |
|
|
|
|
benchmarkMode = "deliver" // deliver/group
|
|
|
|
|
mockUsers = 2000 |
|
|
|
|
eachUserSend = 100 |
|
|
|
|
mockNetUsers []*NetUser |
|
|
|
|
sendCounter int64 = 0 |
|
|
|
|
receiverCounter int64 = 0 |
|
|
|
|
//wsUrls = []string{"ws://192.168.110.36:7001", "ws://192.168.110.36:7003"}
|
|
|
|
|
gatewayHttp = "http://192.168.110.36:7000" |
|
|
|
|
httpClient *http.Client |
|
|
|
|
etcdClient *clientv3.Client |
|
|
|
|
seqId int32 |
|
|
|
|
callbacks map[int32]func(res any) |
|
|
|
|
callbackMutex *sync.Mutex |
|
|
|
|
useMsSum int64 |
|
|
|
|
useMsAvg int64 |
|
|
|
|
httpClient *http.Client |
|
|
|
|
etcdClient *clientv3.Client |
|
|
|
|
seqId int32 |
|
|
|
|
callbacks *collect.ConcurrentMap[int32, func(res any)] |
|
|
|
|
useMsSum int64 |
|
|
|
|
useMsAvg int64 |
|
|
|
|
) |
|
|
|
|
|
|
|
|
|
func init() { |
|
|
|
|
config.InitLogger() |
|
|
|
|
config.InitLogger(false) |
|
|
|
|
|
|
|
|
|
callbacks = make(map[int32]func(res any), 128) |
|
|
|
|
callbackMutex = &sync.Mutex{} |
|
|
|
|
//callbacks = make(map[int32]func(res any), 128)
|
|
|
|
|
//callbackMutex = &sync.Mutex{}
|
|
|
|
|
callbacks = collect.NewConcurrentMap[int32, func(res any)](16, func(k int32) string { |
|
|
|
|
return strconv.Itoa(int(k)) |
|
|
|
|
}) |
|
|
|
|
httpClient = &http.Client{Timeout: 10 * time.Second} |
|
|
|
|
var err error |
|
|
|
|
etcdClient, err = clientv3.New(clientv3.Config{ |
|
|
|
|
@ -68,10 +72,15 @@ func init() {
|
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// main -mode=deliver -users=10 -send=10
|
|
|
|
|
func main() { |
|
|
|
|
flag.StringVar(&benchmarkMode, "mode", "deliver", "benchmark mode: deliver/group") |
|
|
|
|
flag.IntVar(&mockUsers, "users", 10, "mock users") |
|
|
|
|
runtime.GOMAXPROCS(runtime.NumCPU()) |
|
|
|
|
|
|
|
|
|
flag.StringVar(&gatewayHttp, "gateway", "http://192.168.110.41:7000", "http gateway address") // 124.222.131.236:30830
|
|
|
|
|
flag.StringVar(&benchmarkMode, "mode", "group", "benchmark mode: deliver/group") |
|
|
|
|
flag.IntVar(&mockUsers, "users", 100, "mock users") |
|
|
|
|
flag.IntVar(&eachUserSend, "send", 10, "each user send msg count") |
|
|
|
|
flag.Parse() |
|
|
|
|
|
|
|
|
|
if benchmarkMode == "group" { |
|
|
|
|
benchmarkGroup() |
|
|
|
|
@ -109,29 +118,10 @@ func benchmark() {
|
|
|
|
|
// go prof.StartPprof(":8888")
|
|
|
|
|
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background()) |
|
|
|
|
|
|
|
|
|
go records(ctx) |
|
|
|
|
|
|
|
|
|
mockNetUsers = make([]*NetUser, mockUsers) |
|
|
|
|
// initial uids
|
|
|
|
|
for i := 0; i < mockUsers; i++ { |
|
|
|
|
uid := strconv.Itoa(110000 + i) |
|
|
|
|
token, err := getToken(uid) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
conn, err := getConn(token) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
err = handleConn(ctx, uid, token, conn) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
mockNetUsers[i] = &NetUser{ |
|
|
|
|
Uid: uid, |
|
|
|
|
Conn: conn, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
mockUsersConnect(ctx) |
|
|
|
|
|
|
|
|
|
for i := 0; i < mockUsers; i++ { |
|
|
|
|
netUser := mockNetUsers[i] |
|
|
|
|
@ -145,9 +135,6 @@ func benchmark() {
|
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func benchmarkGroup() { |
|
|
|
|
benchmarkMode = "groupDeliver" |
|
|
|
|
runtime.GOMAXPROCS(runtime.NumCPU()) |
|
|
|
|
|
|
|
|
|
ctx, cancel := context.WithCancel(context.Background()) |
|
|
|
|
shutdown.AddHook(cancel) |
|
|
|
|
|
|
|
|
|
@ -163,76 +150,121 @@ func benchmarkGroup() {
|
|
|
|
|
} |
|
|
|
|
groupId := "9527" |
|
|
|
|
|
|
|
|
|
// initial mock users, join to postal group
|
|
|
|
|
mockNetUsers = make([]*NetUser, mockUsers) |
|
|
|
|
for i := 0; i < mockUsers; i++ { |
|
|
|
|
uid := strconv.Itoa(110000 + i) |
|
|
|
|
token, err := getToken(uid) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
conn, err := getConn(token) |
|
|
|
|
mockUsersConnect(ctx) |
|
|
|
|
|
|
|
|
|
// create room
|
|
|
|
|
channel, err := send(mockNetUsers[0].Conn, chat.Chat_ServiceDesc.ServiceName, "RoomCreate", &chat.ReqRoomCreate{Rname: "benchmark"}) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
// wait result
|
|
|
|
|
res := <-channel |
|
|
|
|
if err, failed := res.(error); failed { |
|
|
|
|
panic(err) |
|
|
|
|
} else { |
|
|
|
|
groupId = res.(*chat.ResRoomCreate).Room.Rid |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
for _, user := range mockNetUsers { |
|
|
|
|
channel, err := send(user.Conn, chat.Chat_ServiceDesc.ServiceName, "RoomJoin", &chat.ReqRoomJoin{Rid: groupId}) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
err = handleConn(ctx, uid, token, conn) |
|
|
|
|
if err != nil { |
|
|
|
|
// wait result
|
|
|
|
|
if err, failed := (<-channel).(error); failed { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
mockNetUsers[i] = &NetUser{ |
|
|
|
|
Uid: uid, |
|
|
|
|
Conn: conn, |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// join group
|
|
|
|
|
err = groupDeliver.GroupJoin(context.Background(), uid, []string{groupId}) |
|
|
|
|
shutdown.AddHook(func() { |
|
|
|
|
// delete room
|
|
|
|
|
channel, err := send(mockNetUsers[0].Conn, chat.Chat_ServiceDesc.ServiceName, "RoomDissolve", &chat.ReqRoomDissolve{Rid: groupId}) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
shutdown.AddHook(func() { |
|
|
|
|
if err := groupDeliver.GroupDissolve(context.Background(), groupId); err != nil { |
|
|
|
|
logger.Error("dissolve group error: ", err) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
// wait result
|
|
|
|
|
if err, failed := (<-channel).(error); failed { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
logger.Info("test group dissolved") |
|
|
|
|
}) |
|
|
|
|
|
|
|
|
|
go records(ctx) |
|
|
|
|
|
|
|
|
|
// send group message
|
|
|
|
|
message := &chat.ChatMessage{Sender: "100001", Content: "hello"} |
|
|
|
|
args := &chat.ReqRoomSend{Rid: groupId, Message: "hi"} |
|
|
|
|
for i := 0; i < eachUserSend; i++ { |
|
|
|
|
err := groupDeliver.DeliverGroup(context.Background(), groupId, message) |
|
|
|
|
channel, err := send(mockNetUsers[0].Conn, chat.Chat_ServiceDesc.ServiceName, "RoomSend", args) |
|
|
|
|
if err != nil { |
|
|
|
|
logger.Error("deliver group error:", err) |
|
|
|
|
logger.Error("room send error: ", err) |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
// wait result
|
|
|
|
|
if err, failed := (<-channel).(error); failed { |
|
|
|
|
logger.Error("room send failed: ", err) |
|
|
|
|
} |
|
|
|
|
atomic.AddInt64(&sendCounter, 1) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
shutdown.Await() |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// mockUsersConnect 并行快速创建链接
|
|
|
|
|
func mockUsersConnect(ctx context.Context) { |
|
|
|
|
// initial mock users, join to postal group
|
|
|
|
|
mockNetUsers = make([]*NetUser, mockUsers) |
|
|
|
|
|
|
|
|
|
concurrent := 100 |
|
|
|
|
wg := &sync.WaitGroup{} |
|
|
|
|
wg.Add(concurrent) |
|
|
|
|
each := (mockUsers / concurrent) + 1 |
|
|
|
|
|
|
|
|
|
for i := 0; i < concurrent; i++ { |
|
|
|
|
//begin, end := i*each, (i+1)*each
|
|
|
|
|
//fmt.Printf("%d: %d~%d \n", i, begin, end)
|
|
|
|
|
|
|
|
|
|
go func(segment int) { |
|
|
|
|
defer wg.Done() |
|
|
|
|
|
|
|
|
|
begin, end := segment*each, (segment+1)*each |
|
|
|
|
for i := begin; i < end && i < mockUsers; i++ { |
|
|
|
|
uid := strconv.Itoa(110000 + i) |
|
|
|
|
token, err := getToken(uid) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
conn, err := getConn(token) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
err = handleConn(ctx, uid, token, conn) |
|
|
|
|
if err != nil { |
|
|
|
|
panic(err) |
|
|
|
|
} |
|
|
|
|
mockNetUsers[i] = &NetUser{Uid: uid, Conn: conn} |
|
|
|
|
} |
|
|
|
|
}(i) |
|
|
|
|
} |
|
|
|
|
wg.Wait() |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func sendBatchChatMessage(ctx context.Context, uid string, conn *websocket.Conn, count int) { |
|
|
|
|
r := rand.New(rand.NewSource(time.Now().UnixMilli())) |
|
|
|
|
for i := 0; i < count; i++ { |
|
|
|
|
receiverUid := mockNetUsers[r.Intn(mockUsers)].Uid |
|
|
|
|
args := &chat.ReqSend{ |
|
|
|
|
Receiver: receiverUid, |
|
|
|
|
Content: "hello", |
|
|
|
|
Content: "hi", |
|
|
|
|
} |
|
|
|
|
channel, err := send(conn, chat.Chat_ServiceDesc.ServiceName, "send", args) |
|
|
|
|
channel, err := send(conn, chat.Chat_ServiceDesc.ServiceName, "Send", args) |
|
|
|
|
if err != nil { |
|
|
|
|
fmt.Println("send failed:", err) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
res := <-channel // wait result
|
|
|
|
|
err, failed := res.(error) |
|
|
|
|
if failed { |
|
|
|
|
// wait result
|
|
|
|
|
if err, failed := (<-channel).(error); failed { |
|
|
|
|
fmt.Println("send failed:", err) |
|
|
|
|
} else { |
|
|
|
|
// fmt.Println("send ok:", res)
|
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
@ -297,10 +329,15 @@ func handleConn(ctx context.Context, uid string, token string, conn *websocket.C
|
|
|
|
|
if err != nil { |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
res := <-channel |
|
|
|
|
// fmt.Printf("%v\n", res)
|
|
|
|
|
if e, failed := res.(error); failed { |
|
|
|
|
err = e |
|
|
|
|
ch := time.After(time.Second * 5) |
|
|
|
|
|
|
|
|
|
select { |
|
|
|
|
case <-ch: |
|
|
|
|
err = errors.New("req verify timeout") |
|
|
|
|
case res := <-channel: |
|
|
|
|
if e, failed := res.(error); failed { |
|
|
|
|
err = e |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
@ -347,21 +384,6 @@ func listen(ctx context.Context, conn *websocket.Conn) {
|
|
|
|
|
atomic.AddInt64(&receiverCounter, 1) |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
if header.Type == 4 { |
|
|
|
|
callbackMutex.Lock() |
|
|
|
|
callback, ok := callbacks[header.SeqId] |
|
|
|
|
if ok { |
|
|
|
|
delete(callbacks, header.SeqId) |
|
|
|
|
} |
|
|
|
|
callbackMutex.Unlock() |
|
|
|
|
if !ok { |
|
|
|
|
e = fmt.Errorf("callback %d not found", header.SeqId) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
callback(err) |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// notice message
|
|
|
|
|
if header.Type == protocol.TypeNotice { |
|
|
|
|
switch header.Target { |
|
|
|
|
@ -377,46 +399,55 @@ func listen(ctx context.Context, conn *websocket.Conn) {
|
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
callback, ok := callbacks.LoadAndDelete(header.SeqId) |
|
|
|
|
if !ok { |
|
|
|
|
logger.Errorf("callback %d not found", header.SeqId) |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
if header.Type == protocol.TypeError { |
|
|
|
|
err = errors.New(string(payload.Body)) |
|
|
|
|
callback(err) |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
// rpc response
|
|
|
|
|
|
|
|
|
|
var msg proto.Message |
|
|
|
|
switch header.Svc { |
|
|
|
|
case auth.Auth_ServiceDesc.ServiceName: |
|
|
|
|
switch header.Target { |
|
|
|
|
case "Subject": |
|
|
|
|
// fmt.Println("res verify...")
|
|
|
|
|
msg = &auth.Subject{} |
|
|
|
|
} |
|
|
|
|
case chat.Chat_ServiceDesc.ServiceName: |
|
|
|
|
switch header.Target { |
|
|
|
|
case "ResSend": |
|
|
|
|
// fmt.Println("res send...")
|
|
|
|
|
msg = &chat.ResSend{} |
|
|
|
|
if svc, ok := protoStructs[header.Svc]; ok { |
|
|
|
|
if s, ok := svc[header.Target]; ok { |
|
|
|
|
val := reflect.New(reflect.TypeOf(s).Elem()) |
|
|
|
|
msg = val.Interface().(proto.Message) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
if msg == nil { |
|
|
|
|
logger.Warning("unknown rpc response: ", len(payload.Body), header) |
|
|
|
|
|
|
|
|
|
if msg == nil { // protobuf.Empty
|
|
|
|
|
// logger.Warning("unknown rpc response: ", len(payload.Body), header)
|
|
|
|
|
callback(nil) |
|
|
|
|
continue |
|
|
|
|
} |
|
|
|
|
err = proto.Unmarshal(payload.Body, msg) |
|
|
|
|
if e != nil { |
|
|
|
|
if err != nil { |
|
|
|
|
logger.Error("unmarshal rcp res body error: ", payload, err) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
callbackMutex.Lock() |
|
|
|
|
callback, ok := callbacks[header.SeqId] |
|
|
|
|
if ok { |
|
|
|
|
delete(callbacks, header.SeqId) |
|
|
|
|
} |
|
|
|
|
callbackMutex.Unlock() |
|
|
|
|
if !ok { |
|
|
|
|
e = fmt.Errorf("callback %d not found", header.SeqId) |
|
|
|
|
callback(fmt.Errorf("unmarshal rcp res body error: %v", err)) |
|
|
|
|
return |
|
|
|
|
} |
|
|
|
|
callback(msg) |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
var protoStructs = map[string]map[string]proto.Message{ |
|
|
|
|
auth.Auth_ServiceDesc.ServiceName: { |
|
|
|
|
"Subject": &auth.Subject{}, |
|
|
|
|
}, |
|
|
|
|
chat.Chat_ServiceDesc.ServiceName: { |
|
|
|
|
"ResSend": &chat.ResSend{}, |
|
|
|
|
"ResRoomCreate": &chat.ResRoomCreate{}, |
|
|
|
|
"ResRoomInfo": &chat.ResRoomInfo{}, |
|
|
|
|
"ResRoomList": &chat.ResRoomList{}, |
|
|
|
|
}, |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
func send(conn *websocket.Conn, svc, method string, reqArgs proto.Message) (channel chan any, err error) { |
|
|
|
|
seq := atomic.AddInt32(&seqId, 1) |
|
|
|
|
header := &protocol.Header{ |
|
|
|
|
@ -443,16 +474,14 @@ func send(conn *websocket.Conn, svc, method string, reqArgs proto.Message) (chan
|
|
|
|
|
channel = make(chan any) |
|
|
|
|
|
|
|
|
|
// ready callback
|
|
|
|
|
callbackMutex.Lock() |
|
|
|
|
callbacks[seq] = func(res any) { |
|
|
|
|
callbacks.Store(seq, func(res any) { |
|
|
|
|
ms := time.Now().UnixMilli() - begin |
|
|
|
|
msSum := atomic.AddInt64(&useMsSum, ms) |
|
|
|
|
counter := atomic.AddInt64(&receiverCounter, 1) |
|
|
|
|
atomic.StoreInt64(&useMsAvg, msSum/counter) |
|
|
|
|
|
|
|
|
|
channel <- res |
|
|
|
|
} |
|
|
|
|
callbackMutex.Unlock() |
|
|
|
|
}) |
|
|
|
|
|
|
|
|
|
// send message
|
|
|
|
|
err = conn.WriteMessage(websocket.BinaryMessage, bytes) |
|
|
|
|
|