You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 

150 lines
3.5 KiB

package main
import (
"context"
"io"
"log"
"sig-pub/api/pb"
"sig-pub/pkg/types"
"sig-pub/pkg/zlog"
"time"
"google.golang.org/grpc"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/keepalive"
)
func testSubscribe() {
// 配置 keepalive 参数
keepaliveParams := keepalive.ClientParameters{
Time: 30 * time.Second, // 发送 ping 的间隔
Timeout: 3 * time.Second, // ping 超时
PermitWithoutStream: true, // 允许在没有流的情况下发送 ping
}
conn, err := grpc.NewClient("localhost:8888",
// grpc.WithInsecure(),
grpc.WithKeepaliveParams(keepaliveParams),
grpc.WithTransportCredentials(insecure.NewCredentials()),
)
if err != nil {
log.Fatalf("Failed to connect: %v", err)
}
defer conn.Close()
client := pb.NewExchangeServiceClient(conn)
// subscribeEvents(client)
for {
subscribeEvents(client)
log.Println("Disconnected, retrying in 5 seconds...")
time.Sleep(5 * time.Second) // 重连间隔
}
}
func subscribeEvents(client pb.ExchangeServiceClient) {
// ctx, cancel := context.WithCancel(context.Background())
ctx := context.Background()
stream, err := client.SubscribeKline(ctx)
if err != nil {
return
}
go func() {
<-time.After(time.Second * 60)
// cancel()
err := stream.CloseSend() // 触发 server 端 err = io.EOF
zlog.Infof("连接取消... %v", err)
}()
// 接收消息的goroutine
go func() {
for {
msg, err := stream.Recv()
if err == io.EOF {
zlog.Errorf("服务端关闭连接")
return
}
if err != nil {
zlog.Errorf("接收错误: %v", err)
return
}
for _, k := range msg.Kline.Klines {
kline := new(types.Kline)
kline.ParsePBKline(msg.Kline.Exchange, k)
zlog.Infof("收到服务端消息: %v, %s, %#v", msg.Kline.Exchange, msg.Kline.InstId, kline)
}
}
}()
// 发送消息的goroutine
instIds := []string{"BTC-USDT", "DOGE-USDT-SWAP"}
msg := &pb.ReqStreamSubscribeKline{
SubType: pb.SubscribeType_SUB,
Exchanges: []pb.ExchangeType{pb.ExchangeType_OKX},
InstIds: instIds,
Intervals: []string{
string(types.Interval1s),
string(types.Interval1m),
},
OnlyConfirm: true,
}
zlog.Infof("send stream msg: %#v", msg)
if err = stream.Send(msg); err != nil {
zlog.Errorf("发送失败: %v", err)
return
}
go func() {
<-time.After(10 * time.Second)
msg := &pb.ReqStreamSubscribeKline{
SubType: pb.SubscribeType_UNSUB,
Exchanges: []pb.ExchangeType{pb.ExchangeType_OKX},
InstIds: []string{"BTC-USDT"},
Intervals: []string{
string(types.Interval1s),
string(types.Interval1m),
},
OnlyConfirm: true,
}
if err = stream.Send(msg); err != nil {
zlog.Errorf("发送失败: %v", err)
return
}
zlog.Infof("unsubscribe...")
}()
// go func() {
// scanner := bufio.NewScanner(os.Stdin)
// for scanner.Scan() {
// text := scanner.Text()
// if text == "exit" {
// // 关闭发送端
// if err := stream.CloseSend(); err != nil {
// log.Printf("关闭发送失败: %v", err)
// }
// return
// }
// }
// }()
// 保持连接
<-stream.Context().Done()
log.Println("连接关闭")
// log.Println("subscribeEvents...")
// stream, err := client.SubscribeKline(ctx, &pb.ReqSubscribeKline{Topic: topic})
// if err != nil {
// log.Fatalf("Failed to subscribe: %v", err)
// return
// }
// for {
// event, err := stream.Recv()
// if err != nil {
// log.Printf("Failed to receive event: %v", err)
// break
// }
// fmt.Printf("Received event: InstId=%s, Ts=%v\n", event.InstId, event.Exchange)
// }
}