package main import ( "context" "io" "log" "os" "time" mengyav1 "mengyamonitor-backend/gen/mengyav1" "google.golang.org/grpc" "google.golang.org/grpc/credentials/insecure" ) func runGRPCLoop(addr, serverID, key string) { backoff := time.Second const maxBackoff = 30 * time.Second for { err := runGRPCSession(addr, serverID, key) if err != nil { log.Printf("grpc reporter: session ended: %v (reconnect in %v)", err, backoff) time.Sleep(backoff) if backoff < maxBackoff { backoff *= 2 } continue } backoff = time.Second } } func runGRPCSession(addr, serverID, key string) error { ctx := context.Background() conn, err := grpc.NewClient(addr, grpc.WithTransportCredentials(insecure.NewCredentials())) if err != nil { return err } defer conn.Close() cli := mengyav1.NewAgentIngestClient(conn) stream, err := cli.Connect(ctx) if err != nil { return err } host, _ := os.Hostname() if err := stream.Send(&mengyav1.AgentMessage{ Payload: &mengyav1.AgentMessage_Hello{ Hello: &mengyav1.Hello{ServerId: serverID, AgentKey: key, Hostname: host}, }, }); err != nil { return err } errCh := make(chan error, 1) go func() { for { _, err := stream.Recv() if err != nil { errCh <- err return } } }() tick := time.NewTicker(2 * time.Second) defer tick.Stop() for { select { case err := <-errCh: if err == io.EOF { return nil } return err case <-tick.C: centralMs := measureCollectorToCentralTCP(addr, 3*time.Second) payload, err := BuildDashboardPayloadJSON(centralMs) if err != nil { log.Printf("grpc reporter: collect: %v", err) continue } if err := stream.Send(&mengyav1.AgentMessage{ Payload: &mengyav1.AgentMessage_Metrics{ Metrics: &mengyav1.MetricsReport{ PayloadJson: string(payload), CollectedAtUnixMs: time.Now().UnixMilli(), }, }, }); err != nil { return err } } } }