iT邦幫忙

0

gRPC 双向流长连接:原理、流程与完整示例

  • 分享至 

  • xImage
  •  

示例项目:exampleProject
技术范围:Protocol Buffers、HTTP/2、gRPC 双向流、注册确认、心跳、会话管理、断线重连、请求—响应关联、流式下载。
说明:本文已对原始工程名称、仓库地址和业务标识进行脱敏。代码用于展示设计思路;接入实际项目时,请按目录结构和配置体系调整。

1. gRPC 是什么

gRPC 是一种基于 HTTP/2 的远程过程调用框架。通信双方先用 Protocol Buffers(protobuf)定义服务、方法和消息结构,再由 protoc 生成客户端桩与服务端接口。业务代码操作的是强类型方法和消息,不需要手工拼装 HTTP 请求或 JSON。

gRPC 的几个核心组成如下:

组成 作用
.proto 文件 定义服务、RPC 方法及消息结构,是通信契约
protoc 与插件 生成 Go 消息类型、客户端桩和服务端接口
HTTP/2 提供多路复用、双向流、头部压缩和长连接能力
Channel/ClientConn 客户端到目标服务的逻辑连接,可承载多个 RPC
Stream 一次流式 RPC;本文的长连接实际是一个持续存在的双向流
Metadata/Status 传递调用元数据和标准化错误状态

1.1 四种 RPC 模式

rpc Unary(Request) returns (Response);                  // 一元 RPC
rpc ServerStream(Request) returns (stream Response);    // 服务端流
rpc ClientStream(stream Request) returns (Response);    // 客户端流
rpc BidiStream(stream Request) returns (stream Response); // 双向流

本文使用双向流。客户端和服务端都可以在同一个 Connect RPC 中持续发送消息,适合注册、心跳、服务端主动下发任务以及客户端异步返回结果。

2. gRPC 的工作原理

2.1 从 .proto 到运行时调用

  1. .proto 中定义 ExampleProjectControl.Connect 双向流方法和消息类型。
  2. 使用 protoc 生成 *.pb.go*_grpc.pb.go
  3. 服务端实现生成的 ExampleProjectControlServer 接口并注册到 grpc.Server
  4. 客户端创建 grpc.ClientConn,通过生成的 Client 调用 Connect
  5. protobuf 消息被序列化为二进制帧,经 HTTP/2 stream 传输。
  6. 双方分别调用 SendRecv 持续读写,直到 context 取消、网络中断或任一端结束 RPC。

2.2 “长连接”具体指什么

需要区分四层状态:

  • TCP/TLS 已连接:只代表底层网络通道可用。
  • HTTP/2/gRPC Channel 可用:代表客户端可以发起 RPC。
  • Connect stream 已建立:代表本次双向流 RPC 已开始。
  • 应用注册已确认:代表服务端验证首条注册消息并返回 RegisterAck(Accepted=true)

因此,Dial 成功不等于业务客户端已经在线。示例只有在收到注册确认后才启动心跳并把状态切换为 registered

2.3 HTTP/2 与并发读写

一个 HTTP/2 连接可以复用多个 stream;本文用一个长期存在的 stream 承载控制消息。gRPC Go 允许一个 goroutine 调用 Send、另一个 goroutine 调用 Recv,但不应让多个 goroutine 同时调用同一方向的 Send。因此:

  • 客户端通过 sendMutex 串行化注册、心跳和业务响应;
  • 服务端由唯一 sendLoop 消费 session 发送队列,形成单写者;
  • 收发循环分别运行,避免阻塞彼此。

3. 本示例的完整工作过程

  1. 建立网络连接:exampleProject 客户端建立 TCP/TLS 与 HTTP/2 连接;若部署了 Ingress/LB,由它转发 gRPC 流量到服务端。
  2. 开启双向流:客户端调用 Connect,开始本次双向流 RPC。
  3. 发送注册消息:客户端将 RegisterRequest 作为首条消息,携带 platformprojecthostName 等信息。
  4. 创建会话:服务端校验字段后调用 SessionManager.Add(session);存在同 key 旧会话时,幂等关闭旧会话并替换为新会话。
  5. 确认注册:服务端返回 RegisterAck(Accepted=true),客户端收到后设置 registered=true
  6. 维持在线状态:客户端立即发送一次 Heartbeat,随后定期发送心跳。
  7. 处理业务消息:服务端通过同一双向流下发 QueryRequestReloadConfig;客户端按请求类型执行处理,查询或下载结果通过 QueryResponseDataChunk 返回。

3.1 建连与注册

客户端完成 Dial 后调用 Connect 创建双向流,并把 RegisterRequest 作为第一条消息。服务端只接受注册消息作为首包;校验成功后创建 session,以规范化后的 platform:hostName 作为唯一键,再返回注册确认。

3.2 心跳与超时清理

注册成功后客户端立即发送一次心跳,随后按服务端下发的间隔定时发送。服务端更新 LastSeenAt;清理器定期扫描,超过超时时间的 session 会被关闭、移除并发布状态事件。

推荐满足:

heartbeat interval << heartbeat timeout

例如心跳间隔 15 秒、超时 45 秒。扫描周期与客户端心跳周期最好使用两个独立配置,不要混用同一个字段。

3.3 服务端主动请求与响应关联

服务端先用 FindLive 获取真实在线 session,把带有唯一 request_id 的请求写入发送队列。客户端执行示例查询或下载后,携带相同 request_id 返回响应。服务端通过 pendingQueriespendingDownloads 将异步响应分发给原始调用方。

3.4 连接替换与安全退出

同一客户端重新连接时,新 session 替换旧 session,并关闭旧 session 的 done。旧 Connect handler 退出时会按对象身份执行条件删除,因此不会误删刚建立的新 session。

session.Send 不被关闭,因为多个 goroutine 可能仍在发送;直接关闭会产生 send on closed channel。生命周期结束只关闭幂等的 done 信号,发送方同时监听 done 和 context。

3.5 断线重连

当前状态 发生的事件 下一状态
启动 开始连接 Dialing(尝试连接)
Dialing Dial 与 Connect 成功 StreamOpen(流已开启)
Dialing 连接失败 Backoff(退避等待)
StreamOpen Register 发送成功 AwaitAck(等待注册确认)
StreamOpen Register 发送失败 Backoff
AwaitAck 收到 Accepted=true 的注册确认 Registered(已注册)
AwaitAck 注册被拒绝、确认超时或 Recv 错误 Backoff
Registered 正常交换心跳和业务消息 保持 Registered
Registered stream 断开 Backoff
Backoff 带随机抖动的等待结束 Dialing,重新连接

注册前连续失败采用指数退避,并加入随机抖动,防止大量实例同时重连。一次连接真正注册成功后再断开,下次重连从初始退避开始。

4. Protocol Buffers 定义

建议保存为 proto/exampleproject/agent/v1/control.proto

syntax = "proto3";

package exampleproject.agent.v1;

option go_package = "example.com/exampleProject/proto/exampleproject/agent/v1;agentv1";

service ExampleProjectControl {
  rpc Connect(stream ClientMessage) returns (stream ServerMessage);
}

message ClientMessage {
  string request_id = 1;
  oneof payload {
    RegisterRequest register = 10;
    Heartbeat heartbeat = 11;
    QueryResponse query_response = 12;
    DataChunk data_chunk = 13;
    Ack ack = 14;
    ErrorResponse error = 15;
  }
}

message ServerMessage {
  string request_id = 1;
  oneof payload {
    RegisterAck register_ack = 10;
    Ping ping = 11;
    QueryRequest query_request = 12;
    TailRequest tail_request = 13;
    ReloadConfig reload_config = 14;
  }
}

message RegisterRequest {
  string platform = 1;
  string project = 2;
  string host_name = 3;
  string client_version = 4;
  repeated string users = 5;
  repeated DataSource data_sources = 6;
  int64 started_at_unix = 7;
}

message RegisterAck {
  bool accepted = 1;
  string message = 2;
  int32 heartbeat_interval_seconds = 3;
}

message Heartbeat {
  string platform = 1;
  string project = 2;
  string host_name = 3;
  string status = 4;
  repeated string users = 5;
  repeated DataSource data_sources = 6;
  int64 reported_at_unix = 7;
}

message DataSource {
  string alias = 1;
  bool enabled = 2;
}

message QueryRequest {
  string alias = 1;
  string relative_path = 2;
  int32 lines = 3;
  bool from_end = 4;
  int32 start_line = 5;
  string mode = 6;
  int64 cursor = 7;
}

message QueryResponse {
  repeated string lines = 1;
  int32 total_lines = 2;
  bool has_more = 3;
  int32 start_line = 4;
  int32 end_line = 5;
  int64 cursor = 6;
  int64 file_size = 7;
  repeated FileInfo files = 8;
}

message FileInfo {
  string name = 1;
  string path = 2;
  int64 size = 3;
  int64 mod_time_unix = 4;
  bool is_dir = 5;
  bool is_data = 6;
  string alias = 7;
  string relative_path = 8;
}

message DataChunk {
  bytes data = 1;
  bool eof = 2;
}

message TailRequest {
  string alias = 1;
  string relative_path = 2;
  int32 lines = 3;
}

message ReloadConfig {}
message Ping { int64 sent_at_unix = 1; }
message Ack { string message = 1; }

message ErrorResponse {
  string code = 1;
  string message = 2;
}

生成 Go 代码:

protoc \
  --go_out=. --go_opt=paths=source_relative \
  --go-grpc_out=. --go-grpc_opt=paths=source_relative \
  proto/exampleproject/agent/v1/control.proto

上面的消息名与后文 Go 示例一致。修改 .proto 后应重新生成 *.pb.go*_grpc.pb.go,不要手工编辑生成文件。

5. 文件职责

文件 核心职责
grpc_server_service.go 服务端 启停 gRPC Server、启动超时清理器、优雅停止
grpc_example_project_server_service.go 服务端 Connect 双向流、注册校验、收发循环、心跳与响应处理
grpc_session_manager_service.go 服务端 session 新增、替换、查询、心跳更新、超时清理和并发保护
grpc_client_proxy_service.go 服务端 将 Web/API 请求代理到在线客户端,并关联异步响应
grpc_session_event_service.go 服务端 发布和订阅 session 状态事件
grpc_session_websocket_service.go 服务端 通过 WebSocket 通知浏览器刷新状态
client.go 客户端 Dial、注册确认、心跳、重连、接收指令并返回结果

6. 关键设计要点

  • 只有收到 RegisterAck.Accepted == true 才算注册成功。
  • stream 的第一条客户端消息必须是 RegisterRequest
  • session 唯一键使用规范化后的 platform:hostName
  • List 返回不可写快照;下发消息必须通过 FindLive 获取真实 session。
  • Close() 只关闭 done,不关闭可能被并发写入的 Send
  • 新连接替换旧连接时,通过对象身份检查防止旧 handler 删除新 session。
  • 每个方向只允许一个并发发送者;客户端用互斥锁,服务端用单一发送循环。
  • request_id 必须全局足够唯一,并在请求完成、超时或取消后清理等待表。

7. 完整示例源码

以下代码保留原笔记中的完整实现,并统一使用脱敏后的 exampleProject 命名。

7.1 grpc_server_service.go

package exampleProject

import (
	"context"
	"errors"
	"fmt"
	"net"
	"time"

	"example.com/exampleProject/config"
	"example.com/exampleProject/global"
	agentv1 "example.com/exampleProject/proto/exampleproject/agent/v1"
	"go.uber.org/zap"
	"google.golang.org/grpc"
)

// GrpcServer 是exampleProject gRPC 服务端后台服务。
type GrpcServer struct {
	config config.Grpc

	manager *GrpcSessionManager
}

func NewGrpcServer(cfg config.Grpc) *GrpcServer {
	return &GrpcServer{
		config:  cfg,
		manager: GetGrpcSessionManager(),
	}
}

// Run 启动 gRPC 服务端,并跟随 ctx 生命周期关闭。
func (s *GrpcServer) Run(ctx context.Context) {
	logger := global.GVA_LOG.Named("exampleproject-grpc-server")

	logger.Info(
		"准备启动exampleProject gRPC 服务",
		zap.String("addr", s.config.ListenAddr()),
		zap.Bool("tls", s.config.TLS),
		zap.Duration("heartbeatTimeout", s.config.HeartbeatTimeout()),
		zap.Duration("sessionCheckInterval", s.config.SessionCheckInterval()),
		zap.Int("maxRecvMessageSizeBytes", s.config.MaxRecvMessageSizeBytes()),
		zap.Int("maxSendMessageSizeBytes", s.config.MaxSendMessageSizeBytes()),
	)

	if s.config.TLS {
		logger.Warn("grpc.tls=true 但第一版服务端尚未接入 TLS,当前仍按明文 gRPC 启动")
	}

	listenAddr := s.config.ListenAddr()
	listener, err := net.Listen("tcp", listenAddr)
	if err != nil {
		logger.Error("gRPC 服务监听失败", zap.String("addr", listenAddr), zap.Error(err))
		return
	}

	server := grpc.NewServer(
		grpc.MaxRecvMsgSize(s.config.MaxRecvMessageSizeBytes()),
		grpc.MaxSendMsgSize(s.config.MaxSendMessageSizeBytes()),
	)

	agentv1.RegisterExampleProjectControlServer(
		server,
		NewExampleProjectControlServer(s.manager, s.config),
	)

	serverErr := make(chan error, 1)
	go func() {
		logger.Info("exampleProject gRPC 服务已启动", zap.String("addr", listenAddr))
		if err := server.Serve(listener); err != nil {
			serverErr <- err
			return
		}
		serverErr <- nil
	}()

	cleanerDone := make(chan struct{})
	go s.runSessionCleaner(ctx, cleanerDone, logger)
	logger.Info("gRPC session 清理器已启动", zap.Duration("interval", s.config.SessionCheckInterval()), zap.Duration("timeout", s.config.HeartbeatTimeout()))

	select {
	case <-ctx.Done():
		logger.Info("收到退出信号,准备关闭 gRPC 服务")

	case err := <-serverErr:
		if err != nil && !errors.Is(err, grpc.ErrServerStopped) {
			logger.Error("gRPC 服务异常退出", zap.Error(err))
		}
	}

	stopped := make(chan struct{})
	go func() {
		server.GracefulStop()
		close(stopped)
	}()

	select {
	case <-stopped:
		logger.Info("gRPC 服务已优雅关闭")

	case <-time.After(5 * time.Second):
		logger.Warn("gRPC 服务优雅关闭超时,强制停止")
		server.Stop()
	}

	select {
	case <-cleanerDone:
	case <-time.After(time.Second):
	}
}

func (s *GrpcServer) runSessionCleaner(ctx context.Context, done chan<- struct{}, logger *zap.Logger) {
	defer close(done)

	interval := s.config.SessionCheckInterval()
	timeout := s.config.HeartbeatTimeout()

	if interval <= 0 {
		interval = 15 * time.Second
	}
	if timeout <= 0 {
		timeout = 45 * time.Second
	}

	ticker := time.NewTicker(interval)
	defer ticker.Stop()

	for {
		select {
		case <-ctx.Done():
			return

		case <-ticker.C:
			activeCountBefore := len(s.manager.List())
			logger.Debug(
				"gRPC session 清理器开始扫描",
				zap.Int("activeSessionCount", activeCountBefore),
				zap.Duration("timeout", timeout),
			)
			removed := s.manager.RemoveStale(timeout)
			if len(removed) == 0 {
				logger.Debug("gRPC session 清理器扫描完成,无超时 session", zap.Int("activeSessionCount", len(s.manager.List())))
			}
			for _, key := range removed {
				logger.Warn(
					"gRPC 客户端心跳超时,移除 session",
					zap.String("sessionKey", key),
					zap.Duration("timeout", timeout),
				)
			}
		}
	}
}

func (s *GrpcServer) String() string {
	return fmt.Sprintf("GrpcServer{%s}", s.config.ListenAddr())
}

7.2 grpc_example_project_server_service.go

package exampleProject

import (
	"errors"
	"fmt"
	"io"
	"strings"
	"time"

	"example.com/exampleProject/config"
	"example.com/exampleProject/global"
	agentv1 "example.com/exampleProject/proto/exampleproject/agent/v1"
	"go.uber.org/zap"
)

// ExampleProjectControlServer 实现客户端 gRPC 双向流接入。
type ExampleProjectControlServer struct {
	agentv1.UnimplementedExampleProjectControlServer

	manager *GrpcSessionManager

	config config.Grpc
}

func NewExampleProjectControlServer(manager *GrpcSessionManager, cfg config.Grpc) *ExampleProjectControlServer {
	if manager == nil {
		manager = GetGrpcSessionManager()
	}

	return &ExampleProjectControlServer{
		manager: manager,
		config:  cfg,
	}
}

// Connect 是客户端主动建立的 gRPC 双向长连接。
// 第一条消息必须是 RegisterRequest。
func (s *ExampleProjectControlServer) Connect(stream agentv1.ExampleProjectControl_ConnectServer) error {
	logger := global.GVA_LOG.Named("exampleproject-grpc-stream")
	logger.Info("新的客户端 gRPC Connect stream 已建立,等待首条注册消息")

	firstMessage, err := stream.Recv()
	if err != nil {
		if errors.Is(err, io.EOF) {
			return nil
		}
		return fmt.Errorf("receive first grpc client message: %w", err)
	}

	logger.Info(
		"收到客户端 gRPC 首条消息",
		zap.String("requestID", firstMessage.GetRequestId()),
		zap.String("payloadType", grpcClientMessagePayloadType(firstMessage)),
	)

	registerPayload, ok := firstMessage.Payload.(*agentv1.ClientMessage_Register)
	if !ok || registerPayload.Register == nil {
		_ = stream.Send(&agentv1.ServerMessage{
			RequestId: firstMessage.GetRequestId(),
			Payload: &agentv1.ServerMessage_RegisterAck{
				RegisterAck: &agentv1.RegisterAck{
					Accepted: false,
					Message:  "first grpc client message must be register",
				},
			},
		})
		return errors.New("first grpc client message must be register")
	}

	session, err := s.buildSession(registerPayload.Register)
	if err != nil {
		_ = stream.Send(&agentv1.ServerMessage{
			RequestId: firstMessage.GetRequestId(),
			Payload: &agentv1.ServerMessage_RegisterAck{
				RegisterAck: &agentv1.RegisterAck{
					Accepted: false,
					Message:  err.Error(),
				},
			},
		})
		return err
	}

	s.manager.Add(session)

	global.GVA_LOG.Info(
		"客户端 gRPC 注册成功",
		zap.String("sessionKey", session.Key),
		zap.String("platform", session.Platform),
		zap.String("project", session.Project),
		zap.String("hostName", session.HostName),
		zap.Strings("users", session.Users),
		zap.Int("userCount", len(session.Users)),
		zap.Int("dataSourceCount", len(session.DataSources)),
		zap.Int("sendQueueCap", cap(session.Send)),
	)

	defer func() {
		removed := s.manager.RemoveSession(session)
		if removed {
			global.GVA_LOG.Info(
				"客户端 gRPC session 已移除",
				zap.String("sessionKey", session.Key),
			)
			return
		}

		// 如果同 key 的旧 gRPC stream 被新连接替换,旧 stream 退出时不能删除新 session。
		global.GVA_LOG.Info(
			"客户端 gRPC session 清理已忽略:当前 session 已被新连接替换",
			zap.String("sessionKey", session.Key),
		)
	}()

	sendErr := make(chan error, 1)
	go s.sendLoop(stream, session, sendErr)

	logger.Info(
		"准备发送注册确认 RegisterAck 到客户端",
		zap.String("sessionKey", session.Key),
		zap.String("requestID", firstMessage.GetRequestId()),
	)
	if err := session.SendToClient(&agentv1.ServerMessage{
		RequestId: firstMessage.GetRequestId(),
		Payload: &agentv1.ServerMessage_RegisterAck{
			RegisterAck: &agentv1.RegisterAck{
				Accepted:                 true,
				Message:                  "registered",
				HeartbeatIntervalSeconds: int32(s.config.SessionCheckInterval().Seconds()),
			},
		},
	}); err != nil {
		logger.Warn("注册确认写入 session.Send 队列失败", zap.String("sessionKey", session.Key), zap.Error(err))
		return fmt.Errorf("enqueue grpc register ack: %w", err)
	}

	recvErr := make(chan error, 1)
	go func() {
		recvErr <- s.receiveClientMessages(stream, session, logger)
	}()

	select {
	case err := <-sendErr:
		return err
	case err := <-recvErr:
		return err
	case <-session.Done():
		logger.Info("gRPC session 已关闭,Connect handler 退出", zap.String("sessionKey", session.Key))
		return nil
	case <-stream.Context().Done():
		return stream.Context().Err()
	}
}

// receiveClientMessages 单独负责阻塞式 Recv。
// Connect 主 goroutine 可以同时监听 session.Done,确保 session 被替换或超时清理时整条 RPC 及时退出。
func (s *ExampleProjectControlServer) receiveClientMessages(
	stream agentv1.ExampleProjectControl_ConnectServer,
	session *GrpcClientSession,
	logger *zap.Logger,
) error {
	for {
		message, err := stream.Recv()
		if err != nil {
			if errors.Is(err, io.EOF) {
				return nil
			}
			return fmt.Errorf("receive grpc client message: %w", err)
		}

		if message == nil {
			logger.Warn("收到空的客户端 gRPC 消息", zap.String("sessionKey", session.Key))
			continue
		}

		logger.Debug(
			"收到客户端 gRPC 消息",
			zap.String("sessionKey", session.Key),
			zap.String("requestID", message.GetRequestId()),
			zap.String("payloadType", grpcClientMessagePayloadType(message)),
		)

		if err := s.handleClientMessage(session, message); err != nil {
			global.GVA_LOG.Warn(
				"处理客户端 gRPC 消息失败",
				zap.String("sessionKey", session.Key),
				zap.Error(err),
			)
		}
	}
}

func (s *ExampleProjectControlServer) sendLoop(
	stream agentv1.ExampleProjectControl_ConnectServer,
	session *GrpcClientSession,
	errCh chan<- error,
) {
	logger := global.GVA_LOG.Named("exampleproject-grpc-stream")
	logger.Info(
		"gRPC 服务端发送循环已启动",
		zap.String("sessionKey", session.Key),
		zap.Int("sendQueueCap", cap(session.Send)),
	)

	for {
		select {
		case <-session.Done():
			logger.Info("gRPC session 已关闭,发送循环退出", zap.String("sessionKey", session.Key))
			return

		case <-stream.Context().Done():
			err := stream.Context().Err()
			logger.Warn("gRPC stream context 已结束,发送循环退出", zap.String("sessionKey", session.Key), zap.Error(err))
			select {
			case errCh <- err:
			default:
			}
			return

		case message := <-session.Send:
			if message == nil {
				logger.Warn("从 session.Send 队列取到空消息", zap.String("sessionKey", session.Key))
				continue
			}

			logger.Info(
				"准备通过 gRPC stream.Send 下发消息给客户端",
				zap.String("sessionKey", session.Key),
				zap.String("requestID", message.GetRequestId()),
				zap.String("payloadType", grpcServerMessagePayloadType(message)),
			)

			if err := stream.Send(message); err != nil {
				logger.Error(
					"gRPC stream.Send 下发消息失败",
					zap.String("sessionKey", session.Key),
					zap.String("requestID", message.GetRequestId()),
					zap.String("payloadType", grpcServerMessagePayloadType(message)),
					zap.Error(err),
				)
				select {
				case errCh <- fmt.Errorf("send grpc server message: %w", err):
				default:
				}
				return
			}

			logger.Info(
				"gRPC stream.Send 下发消息成功",
				zap.String("sessionKey", session.Key),
				zap.String("requestID", message.GetRequestId()),
				zap.String("payloadType", grpcServerMessagePayloadType(message)),
			)
		}
	}
}

func (s *ExampleProjectControlServer) buildSession(register *agentv1.RegisterRequest) (*GrpcClientSession, error) {
	if register == nil {
		return nil, errors.New("register request is nil")
	}

	platform := strings.TrimSpace(register.GetPlatform())
	hostName := strings.TrimSpace(register.GetHostName())
	// project 是子项目名称,用于服务端/前端项目树渲染。
	// project 不要求等于 data_sources[].alias。
	project := strings.TrimSpace(register.GetProject())

	if platform == "" {
		return nil, errors.New("register.platform is empty")
	}
	if hostName == "" {
		return nil, errors.New("register.host_name is empty")
	}
	if project == "" {
		return nil, errors.New("register.project is empty")
	}

	now := time.Now().UTC()
	startedAt := now
	if register.GetStartedAtUnix() > 0 {
		startedAt = time.Unix(register.GetStartedAtUnix(), 0).UTC()
	}

	dataSources := normalizeDataSourcesFromProto(register.GetDataSources())
	if len(dataSources) == 0 {
		return nil, errors.New("register.data_sources is empty: data_sources[].alias must come from example-project.data-dirs[].alias")
	}

	users := cleanStringList(register.GetUsers())
	if len(users) == 0 {
		return nil, errors.New("register.users is empty")
	}

	key := BuildClientSessionKey(platform, hostName)

	return &GrpcClientSession{
		Key:           key,
		Platform:      platform,
		Project:       project,
		HostName:      hostName,
		ClientVersion: strings.TrimSpace(register.GetClientVersion()),
		Users:         users,
		DataSources:    dataSources,
		StartedAt:     startedAt,
		RegisteredAt:  now,
		LastSeenAt:    now,
		Status:        "online",
		Send:          make(chan *agentv1.ServerMessage, 64),
		done:          make(chan struct{}),
	}, nil
}

func cleanStringList(values []string) []string {
	if len(values) == 0 {
		return []string{}
	}

	seen := make(map[string]struct{}, len(values))
	result := make([]string, 0, len(values))

	for _, raw := range values {
		value := strings.TrimSpace(raw)
		if value == "" {
			continue
		}

		if _, exists := seen[value]; exists {
			continue
		}

		seen[value] = struct{}{}
		result = append(result, value)
	}

	return result
}

func (s *ExampleProjectControlServer) handleClientMessage(session *GrpcClientSession, message *agentv1.ClientMessage) error {
	logger := global.GVA_LOG.Named("exampleproject-grpc-stream")
	logger.Debug(
		"开始处理客户端 gRPC 消息",
		zap.String("sessionKey", session.Key),
		zap.String("requestID", message.GetRequestId()),
		zap.String("payloadType", grpcClientMessagePayloadType(message)),
	)

	switch payload := message.Payload.(type) {
	case *agentv1.ClientMessage_Heartbeat:
		err := s.handleHeartbeat(session, payload.Heartbeat)
		if err == nil {
			logger.Debug(
				"客户端心跳处理成功",
				zap.String("sessionKey", session.Key),
				zap.String("status", payload.Heartbeat.GetStatus()),
				zap.Int64("reportedAtUnix", payload.Heartbeat.GetReportedAtUnix()),
			)
		}
		return err

	case *agentv1.ClientMessage_Ack:
		global.GVA_LOG.Debug(
			"收到客户端 ACK",
			zap.String("sessionKey", session.Key),
			zap.String("requestID", message.GetRequestId()),
			zap.String("message", payload.Ack.GetMessage()),
		)
		return nil

	case *agentv1.ClientMessage_Error:
		if HandleGrpcClientMessageResult(message.GetRequestId(), payload.Error) {
			logger.Info(
				"客户端错误响应已分发给等待中的请求",
				zap.String("sessionKey", session.Key),
				zap.String("requestID", message.GetRequestId()),
				zap.String("code", payload.Error.GetCode()),
			)
			return nil
		}
		global.GVA_LOG.Warn(
			"收到客户端错误消息",
			zap.String("sessionKey", session.Key),
			zap.String("requestID", message.GetRequestId()),
			zap.String("code", payload.Error.GetCode()),
			zap.String("message", payload.Error.GetMessage()),
		)
		return nil

	case *agentv1.ClientMessage_QueryResponse:
		if HandleGrpcClientMessageResult(message.GetRequestId(), payload.QueryResponse) {
			logger.Info(
				"客户端数据查询响应已分发给等待中的请求",
				zap.String("sessionKey", session.Key),
				zap.String("requestID", message.GetRequestId()),
				zap.Int("lineCount", len(payload.QueryResponse.GetLines())),
			)
			return nil
		}
		global.GVA_LOG.Debug(
			"收到客户端数据查询响应,但没有等待者",
			zap.String("sessionKey", session.Key),
			zap.String("requestID", message.GetRequestId()),
			zap.Int("lineCount", len(payload.QueryResponse.GetLines())),
		)
		return nil

	case *agentv1.ClientMessage_DataChunk:
		if HandleGrpcClientMessageResult(message.GetRequestId(), payload.DataChunk) {
			logger.Debug(
				"客户端数据下载分块已分发给等待中的请求",
				zap.String("sessionKey", session.Key),
				zap.String("requestID", message.GetRequestId()),
				zap.Int("bytes", len(payload.DataChunk.GetData())),
				zap.Bool("eof", payload.DataChunk.GetEof()),
			)
			return nil
		}
		global.GVA_LOG.Debug(
			"收到客户端数据下载分块,但没有等待者",
			zap.String("sessionKey", session.Key),
			zap.String("requestID", message.GetRequestId()),
			zap.Int("bytes", len(payload.DataChunk.GetData())),
			zap.Bool("eof", payload.DataChunk.GetEof()),
		)
		return nil

	case *agentv1.ClientMessage_Register:
		return errors.New("duplicate register message")

	default:
		return nil
	}
}

func (s *ExampleProjectControlServer) handleHeartbeat(session *GrpcClientSession, heartbeat *agentv1.Heartbeat) error {
	if heartbeat == nil {
		return errors.New("heartbeat is nil")
	}

	platform := strings.TrimSpace(heartbeat.GetPlatform())
	hostName := strings.TrimSpace(heartbeat.GetHostName())
	project := strings.TrimSpace(heartbeat.GetProject())
	// heartbeat.project 是子项目名称,用于服务端/前端项目树渲染;不要求等于 data_sources[].alias。
	// 客户端数据目录变更时,应通过 heartbeat.data_sources[].alias 刷新服务端数据查询 alias。
	if platform != "" && platform != session.Platform {
		return fmt.Errorf("heartbeat platform mismatch: %s != %s", platform, session.Platform)
	}
	if hostName != "" && hostName != session.HostName {
		return fmt.Errorf("heartbeat hostName mismatch: %s != %s", hostName, session.HostName)
	}
	if project != "" && strings.TrimSpace(session.Project) != "" && project != session.Project {
		return fmt.Errorf("heartbeat project mismatch: %s != %s", project, session.Project)
	}

	users := cleanStringList(heartbeat.GetUsers())
	if len(users) == 0 {
		return errors.New("heartbeat.users is empty")
	}

	reportedAt := time.Now().UTC()
	if heartbeat.GetReportedAtUnix() > 0 {
		reportedAt = time.Unix(heartbeat.GetReportedAtUnix(), 0).UTC()
	}

	dataSources := normalizeDataSourcesFromProto(heartbeat.GetDataSources())

	return s.manager.UpdateHeartbeat(session.Key, project, heartbeat.GetStatus(), reportedAt, users, dataSources)
}

func normalizeDataSourcesFromProto(values []*agentv1.DataSource) []GrpcDataSource {
	if len(values) == 0 {
		return []GrpcDataSource{}
	}

	seenAlias := make(map[string]struct{}, len(values))
	result := make([]GrpcDataSource, 0, len(values))

	for _, source := range values {
		if source == nil {
			continue
		}

		alias := strings.TrimSpace(source.GetAlias())
		if alias == "" || !source.GetEnabled() {
			continue
		}

		if _, exists := seenAlias[alias]; exists {
			continue
		}
		seenAlias[alias] = struct{}{}

		result = append(result, GrpcDataSource{
			Alias:   alias,
			Enabled: true,
		})
	}

	return result
}

func grpcClientMessagePayloadType(message *agentv1.ClientMessage) string {
	if message == nil {
		return "<nil>"
	}
	switch message.Payload.(type) {
	case *agentv1.ClientMessage_Register:
		return "Register"
	case *agentv1.ClientMessage_Heartbeat:
		return "Heartbeat"
	case *agentv1.ClientMessage_QueryResponse:
		return "QueryResponse"
	case *agentv1.ClientMessage_DataChunk:
		return "DataChunk"
	case *agentv1.ClientMessage_Ack:
		return "Ack"
	case *agentv1.ClientMessage_Error:
		return "Error"
	case nil:
		return "<nil>"
	default:
		return fmt.Sprintf("%T", message.Payload)
	}
}

func grpcServerMessagePayloadType(message *agentv1.ServerMessage) string {
	if message == nil {
		return "<nil>"
	}
	switch message.Payload.(type) {
	case *agentv1.ServerMessage_RegisterAck:
		return "RegisterAck"
	case *agentv1.ServerMessage_Ping:
		return "Ping"
	case *agentv1.ServerMessage_QueryRequest:
		return "QueryRequest"
	case *agentv1.ServerMessage_TailRequest:
		return "TailRequest"
	case *agentv1.ServerMessage_ReloadConfig:
		return "ReloadConfig"
	case nil:
		return "<nil>"
	default:
		return fmt.Sprintf("%T", message.Payload)
	}
}

7.3 grpc_session_manager_service.go

package exampleProject

import (
	"context"
	"errors"
	"fmt"
	"sort"
	"strings"
	"sync"
	"time"

	agentv1 "example.com/exampleProject/proto/exampleproject/agent/v1"
)

var defaultGrpcSessionManager = NewGrpcSessionManager()

func GetGrpcSessionManager() *GrpcSessionManager {
	return defaultGrpcSessionManager
}

type GrpcSessionManager struct {
	mu sync.RWMutex

	sessions map[string]*GrpcClientSession
}

func NewGrpcSessionManager() *GrpcSessionManager {
	return &GrpcSessionManager{
		sessions: make(map[string]*GrpcClientSession),
	}
}

// GrpcClientSession 表示一个通过 gRPC Connect 注册成功的客户端长连接。
type GrpcClientSession struct {
	Key string

	Platform string

	// Project 是子项目名称,来自客户端 config.yaml 中 example-project.project。
	// 它用于服务端/前端项目树渲染;不要求等于 DataSources[].Alias。
	Project string

	HostName string

	ClientVersion string

	Users []string

	DataSources []GrpcDataSource

	StartedAt time.Time

	RegisteredAt time.Time

	LastSeenAt time.Time

	Status string

	Send chan *agentv1.ServerMessage

	done chan struct{}

	closeOnce sync.Once
}

type GrpcDataSource struct {
	Alias string `json:"alias"`

	Enabled bool `json:"enabled"`
}

func (s *GrpcClientSession) Close() {
	if s == nil {
		return
	}

	s.closeOnce.Do(func() {
		close(s.done)
	})
}

func (s *GrpcClientSession) Done() <-chan struct{} {
	if s == nil {
		closed := make(chan struct{})
		close(closed)
		return closed
	}
	return s.done
}

func (s *GrpcClientSession) SendToClient(message *agentv1.ServerMessage) error {
	return s.SendToClientContext(context.Background(), message)
}

func (s *GrpcClientSession) SendToClientContext(ctx context.Context, message *agentv1.ServerMessage) error {
	if s == nil {
		return errors.New("client session is nil")
	}
	if message == nil {
		return errors.New("server message is nil")
	}
	if ctx == nil {
		ctx = context.Background()
	}
	if s.Send == nil {
		return errors.New("client session send channel is nil: probably using a snapshot from List(), use FindLive instead")
	}
	if s.done == nil {
		return errors.New("client session done channel is nil: probably using a snapshot from List(), use FindLive instead")
	}

	select {
	case <-s.done:
		return errors.New("client session is closed")

	case <-ctx.Done():
		return ctx.Err()

	case s.Send <- message:
		return nil
	}
}

func BuildClientSessionKey(platform string, hostName string) string {
	platform = sanitizeSessionKeySegment(platform)
	hostName = sanitizeSessionKeySegment(hostName)
	return strings.Join([]string{platform, hostName}, ":")
}

func sanitizeSessionKeySegment(value string) string {
	value = strings.TrimSpace(value)
	if value == "" {
		return "unknown"
	}

	replacer := strings.NewReplacer(
		":", "_",
		"/", "_",
		"\\", "_",
		" ", "_",
		"\t", "_",
		"\n", "_",
		"\r", "_",
	)

	return replacer.Replace(value)
}

func (m *GrpcSessionManager) Add(session *GrpcClientSession) {
	if m == nil || session == nil {
		return
	}

	var replaced *GrpcClientSession

	m.mu.Lock()
	if old, exists := m.sessions[session.Key]; exists && old != nil && old != session {
		replaced = cloneSessionForEvent(old)
		old.Close()
	}

	if strings.TrimSpace(session.Status) == "" {
		session.Status = "online"
	}
	if session.RegisteredAt.IsZero() {
		session.RegisteredAt = time.Now().UTC()
	}
	if session.LastSeenAt.IsZero() {
		session.LastSeenAt = time.Now().UTC()
	}

	m.sessions[session.Key] = session
	online := cloneSessionForEvent(session)
	m.mu.Unlock()

	if replaced != nil {
		PublishGrpcSessionEventFromSession(GrpcSessionEventClientReplace, replaced)
	}
	PublishGrpcSessionEventFromSession(GrpcSessionEventClientOnline, online)
}

func (m *GrpcSessionManager) Remove(key string) {
	removed, ok := m.removeByKey(key, nil, false)
	if ok && removed != nil {
		PublishGrpcSessionEventFromSession(GrpcSessionEventClientOffline, removed)
	}
}

// RemoveSession 只在 manager 当前保存的 session 与传入 session 是同一个对象时才删除。
//
// 这个方法用于 gRPC Connect defer 清理,避免“旧连接被新连接替换后,旧连接退出时把新 session 删掉”。
func (m *GrpcSessionManager) RemoveSession(session *GrpcClientSession) bool {
	if session == nil {
		return false
	}

	removed, ok := m.removeByKey(session.Key, session, true)
	if ok && removed != nil {
		PublishGrpcSessionEventFromSession(GrpcSessionEventClientOffline, removed)
	}
	return ok
}

func (m *GrpcSessionManager) removeByKey(key string, expected *GrpcClientSession, requireSame bool) (*GrpcClientSession, bool) {
	if m == nil {
		return nil, false
	}

	key = strings.TrimSpace(key)
	if key == "" {
		return nil, false
	}

	var removed *GrpcClientSession

	m.mu.Lock()
	if session, exists := m.sessions[key]; exists && session != nil {
		if requireSame && session != expected {
			m.mu.Unlock()
			return nil, false
		}

		removed = cloneSessionForEvent(session)
		removed.Status = "offline"
		session.Close()
		delete(m.sessions, key)
	}
	m.mu.Unlock()

	return removed, removed != nil
}

func (m *GrpcSessionManager) Get(key string) (*GrpcClientSession, bool) {
	if m == nil {
		return nil, false
	}

	key = strings.TrimSpace(key)
	if key == "" {
		return nil, false
	}

	m.mu.RLock()
	defer m.mu.RUnlock()

	session, ok := m.sessions[key]
	return session, ok
}

// FindLive 根据外部传入的标识查找 manager 内部保存的真实 session。
// 支持传入 sessionKey、hostName、platform:hostName、clientId 等候选值。
func (m *GrpcSessionManager) FindLive(
	identifier string,
	candidates func(*GrpcClientSession) []string,
) (*GrpcClientSession, bool) {
	if m == nil {
		return nil, false
	}

	identifier = strings.TrimSpace(identifier)
	if identifier == "" {
		return nil, false
	}

	m.mu.RLock()
	defer m.mu.RUnlock()

	if session, ok := m.sessions[identifier]; ok && session != nil {
		return session, true
	}

	for _, session := range m.sessions {
		if session == nil {
			continue
		}

		baseCandidates := []string{
			session.Key,
			session.HostName,
			BuildClientSessionKey(session.Platform, session.HostName),
			strings.Join([]string{session.Platform, session.HostName}, "-"),
		}

		if candidates != nil {
			baseCandidates = append(baseCandidates, candidates(session)...)
		}

		for _, candidate := range baseCandidates {
			if strings.TrimSpace(candidate) == identifier {
				return session, true
			}
		}
	}

	return nil, false
}

func (m *GrpcSessionManager) List() []*GrpcClientSession {
	if m == nil {
		return []*GrpcClientSession{}
	}

	m.mu.RLock()
	defer m.mu.RUnlock()

	list := make([]*GrpcClientSession, 0, len(m.sessions))
	for _, session := range m.sessions {
		if session == nil {
			continue
		}
		list = append(list, cloneSession(session))
	}

	sort.Slice(list, func(i, j int) bool {
		return list[i].Key < list[j].Key
	})

	return list
}

func cloneSession(session *GrpcClientSession) *GrpcClientSession {
	if session == nil {
		return nil
	}

	clone := *session
	clone.Users = append([]string(nil), session.Users...)
	clone.DataSources = append([]GrpcDataSource(nil), session.DataSources...)
	clone.Send = nil
	clone.done = nil
	return &clone
}

func cloneSessionForEvent(session *GrpcClientSession) *GrpcClientSession {
	return cloneSession(session)
}

func (m *GrpcSessionManager) UpdateHeartbeat(
	key string,
	project string,
	status string,
	heartbeatAt time.Time,
	users []string,
	dataSources []GrpcDataSource,
) error {
	if m == nil {
		return errors.New("session manager is nil")
	}

	key = strings.TrimSpace(key)
	if key == "" {
		return errors.New("session key is empty")
	}

	project = strings.TrimSpace(project)
	cleanUsers := normalizeSessionUsers(users)
	if len(cleanUsers) == 0 {
		return errors.New("heartbeat.users is empty")
	}
	cleanDataSources := normalizeSessionDataSources(dataSources)

	var changed *GrpcClientSession

	m.mu.Lock()
	session, exists := m.sessions[key]
	if !exists || session == nil {
		m.mu.Unlock()
		return fmt.Errorf("session not found: %s", key)
	}

	if len(session.Users) == 0 {
		session.Users = cleanUsers
		changed = cloneSessionForEvent(session)
	} else if !sameStringSet(session.Users, cleanUsers) {
		m.mu.Unlock()
		return fmt.Errorf(
			"heartbeat users mismatch: register=%v heartbeat=%v",
			session.Users,
			cleanUsers,
		)
	}

	if project != "" {
		if strings.TrimSpace(session.Project) == "" {
			session.Project = project
			changed = cloneSessionForEvent(session)
		} else if session.Project != project {
			m.mu.Unlock()
			return fmt.Errorf("heartbeat project mismatch: register=%s heartbeat=%s", session.Project, project)
		}
	}

	if heartbeatAt.IsZero() {
		heartbeatAt = time.Now().UTC()
	}

	status = strings.TrimSpace(status)
	if status == "" {
		status = "online"
	}

	if len(cleanDataSources) > 0 && !sameDataSourceSet(session.DataSources, cleanDataSources) {
		// 客户端 example-project.data-dirs 支持多个;当配置 reload 后,心跳可以同步最新 alias 列表。
		// DataSources[].Alias 仅用于数据查询,不要求等于 Project。
		session.DataSources = cleanDataSources
		changed = cloneSessionForEvent(session)
	}

	if session.Status != status {
		session.Status = status
		changed = cloneSessionForEvent(session)
	} else {
		session.Status = status
	}
	session.LastSeenAt = heartbeatAt.UTC()
	if changed != nil {
		changed = cloneSessionForEvent(session)
	}
	m.mu.Unlock()

	if changed != nil {
		PublishGrpcSessionEventFromSession(GrpcSessionEventClientUpdate, changed)
	}
	return nil
}

func normalizeSessionDataSources(values []GrpcDataSource) []GrpcDataSource {
	if len(values) == 0 {
		return []GrpcDataSource{}
	}

	seen := make(map[string]struct{}, len(values))
	result := make([]GrpcDataSource, 0, len(values))

	for _, item := range values {
		alias := strings.TrimSpace(item.Alias)
		if alias == "" || !item.Enabled {
			continue
		}
		if _, exists := seen[alias]; exists {
			continue
		}
		seen[alias] = struct{}{}
		result = append(result, GrpcDataSource{Alias: alias, Enabled: true})
	}

	sort.Slice(result, func(i, j int) bool { return result[i].Alias < result[j].Alias })
	return result
}

func sameDataSourceSet(left []GrpcDataSource, right []GrpcDataSource) bool {
	left = normalizeSessionDataSources(left)
	right = normalizeSessionDataSources(right)
	if len(left) != len(right) {
		return false
	}
	for i := range left {
		if left[i].Alias != right[i].Alias || left[i].Enabled != right[i].Enabled {
			return false
		}
	}
	return true
}

func normalizeSessionUsers(values []string) []string {
	if len(values) == 0 {
		return []string{}
	}

	seen := make(map[string]struct{}, len(values))
	result := make([]string, 0, len(values))

	for _, raw := range values {
		value := strings.TrimSpace(raw)
		if value == "" {
			continue
		}
		if _, exists := seen[value]; exists {
			continue
		}
		seen[value] = struct{}{}
		result = append(result, value)
	}

	sort.Strings(result)
	return result
}

func sameStringSet(left []string, right []string) bool {
	left = normalizeSessionUsers(left)
	right = normalizeSessionUsers(right)

	if len(left) != len(right) {
		return false
	}

	for i := range left {
		if left[i] != right[i] {
			return false
		}
	}

	return true
}

func (m *GrpcSessionManager) RemoveStale(timeout time.Duration) []string {
	if m == nil || timeout <= 0 {
		return nil
	}

	now := time.Now().UTC()
	removed := make([]string, 0)
	removedSessions := make([]*GrpcClientSession, 0)
	removedKeysOnly := make([]string, 0)

	m.mu.Lock()
	for key, session := range m.sessions {
		if session == nil {
			delete(m.sessions, key)
			removed = append(removed, key)
			removedKeysOnly = append(removedKeysOnly, key)
			continue
		}

		if now.Sub(session.LastSeenAt) <= timeout {
			continue
		}

		eventSession := cloneSessionForEvent(session)
		eventSession.Status = "timeout"

		session.Close()
		delete(m.sessions, key)

		removed = append(removed, key)
		removedSessions = append(removedSessions, eventSession)
	}
	m.mu.Unlock()

	for _, session := range removedSessions {
		PublishGrpcSessionEventFromSession(GrpcSessionEventClientTimeout, session)
	}
	for _, key := range removedKeysOnly {
		PublishGrpcSessionEventFromKey(GrpcSessionEventClientTimeout, key)
	}

	return removed
}

7.4 grpc_client_proxy_service.go

package exampleProject

import (
	"context"
	"crypto/rand"
	"encoding/hex"
	"encoding/json"
	"errors"
	"fmt"
	"hash/fnv"
	"io"
	"sort"
	"strings"
	"sync"
	"time"

	"example.com/exampleProject/model/system/request"
	agentv1 "example.com/exampleProject/proto/exampleproject/agent/v1"
)

const (
	defaultGrpcRequestTimeout = 30 * time.Second
	defaultGrpcSendTimeout    = 5 * time.Second
	downloadChunkChannelSize  = 32
)

type GrpcClientProxyService struct {
	manager *GrpcSessionManager
}

func NewGrpcClientProxyService() *GrpcClientProxyService {
	return &GrpcClientProxyService{manager: GetGrpcSessionManager()}
}

// GrpcClientHostView 是返回前端的在线客户端视图。
type GrpcClientHostView struct {
	ID       uint   `json:"ID"`
	Name     string `json:"name"`
	Url      string `json:"url"`
	Comments string `json:"comments"`
	GrpcKey  string `json:"grpcKey"`
	ClientID string `json:"clientId"`
	Platform string `json:"platform"`
  
	Project       string          `json:"project"`
	HostName      string          `json:"hostName"`
	ClientVersion string          `json:"clientVersion"`
	Users         []string        `json:"users"`
	ExampleProjectUsers   []string        `json:"exampleProjectUsers"`
	DataAliases    []string        `json:"dataAliases"`
	DataDirs       []GrpcDataSource `json:"dataDirs"`
	Status        string          `json:"status"`
	RegisteredAt  string          `json:"registeredAt"`
	LastSeenAt    string          `json:"lastSeenAt"`
	HeartbeatAt   string          `json:"heartbeatAt"`
}

type GrpcClientDataDirView struct {
	Alias    string `json:"alias"`
	DataDir   string `json:"data_dir"`
	HostName string `json:"host_name"`
	HostURL  string `json:"host_url"`
	Enabled  bool   `json:"enabled"`
}

type GrpcFileInfoView struct {
	Name    string    `json:"name"`
	Path    string    `json:"path"`
	Size    int64     `json:"size"`
	ModTime time.Time `json:"modTime"`
	IsDir   bool      `json:"isDir"`
	IsData   bool      `json:"isData"`
	Alias   string    `json:"alias,omitempty"`
	Rel     string    `json:"rel,omitempty"`
}

type ReadDataRequest struct {
	HostID    string
	FilePath  string
	Lines     int
	FromEnd   bool
	StartLine int
	Mode      string
	Cursor    int64
}

type ReadDataResult struct {
	Lines      []string `json:"lines"`
	TotalLines int      `json:"totalLines"`
	HasMore    bool     `json:"hasMore"`
	StartLine  int      `json:"startLine"`
	EndLine    int      `json:"endLine"`
	Cursor     int64    `json:"cursor"`
	FileSize   int64    `json:"fileSize"`
}

type DownloadDataRequest struct {
	HostID   string
	FilePath string
	Writer   io.Writer
}

type grpcQueryResult struct {
	response *agentv1.QueryResponse
	err      error
}

type grpcDownloadChunk struct {
	data []byte
	eof  bool
	err  error
}

var (
	pendingQueryMu sync.Mutex
	pendingQueries = make(map[string]chan grpcQueryResult)

	pendingDownloadMu sync.Mutex
	pendingDownloads  = make(map[string]chan grpcDownloadChunk)
)

func HandleGrpcClientMessageResult(requestID string, payload any) bool {
	requestID = strings.TrimSpace(requestID)
	if requestID == "" {
		return false
	}

	switch msg := payload.(type) {
	case *agentv1.QueryResponse:
		pendingQueryMu.Lock()
		ch := pendingQueries[requestID]
		pendingQueryMu.Unlock()
		if ch == nil {
			return false
		}
		select {
		case ch <- grpcQueryResult{response: msg}:
		default:
		}
		return true

	case *agentv1.DataChunk:
		pendingDownloadMu.Lock()
		ch := pendingDownloads[requestID]
		pendingDownloadMu.Unlock()
		if ch == nil {
			return false
		}
		select {
		case ch <- grpcDownloadChunk{data: msg.GetData(), eof: msg.GetEof()}:
		case <-time.After(time.Second):
		}
		return true

	case *agentv1.ErrorResponse:
		err := errors.New(msg.GetMessage())
		if strings.TrimSpace(msg.GetCode()) != "" {
			err = fmt.Errorf("%s: %s", msg.GetCode(), msg.GetMessage())
		}

		pendingQueryMu.Lock()
		queryCh := pendingQueries[requestID]
		pendingQueryMu.Unlock()
		if queryCh != nil {
			select {
			case queryCh <- grpcQueryResult{err: err}:
			default:
			}
			return true
		}

		pendingDownloadMu.Lock()
		downloadCh := pendingDownloads[requestID]
		pendingDownloadMu.Unlock()
		if downloadCh != nil {
			select {
			case downloadCh <- grpcDownloadChunk{err: err}:
			default:
			}
			return true
		}
	}

	return false
}

func (s *GrpcClientProxyService) ListClientHosts(ctx context.Context, pageInfo request.ClientHostsListSearch) ([]GrpcClientHostView, int64, error) {
	_ = ctx
	list := s.allHostViews()
	total := int64(len(list))

	page := pageInfo.Page
	pageSize := pageInfo.PageSize
	if page > 0 && pageSize > 0 {
		start := (page - 1) * pageSize
		if start >= len(list) {
			return []GrpcClientHostView{}, total, nil
		}
		end := start + pageSize
		if end > len(list) {
			end = len(list)
		}
		list = list[start:end]
	}

	return list, total, nil
}

func (s *GrpcClientProxyService) SearchClientHosts(ctx context.Context, keyword string) ([]GrpcClientHostView, error) {
	_ = ctx
	keyword = strings.ToLower(strings.TrimSpace(keyword))
	if keyword == "" {
		return s.allHostViews(), nil
	}

	result := make([]GrpcClientHostView, 0)
	for _, item := range s.allHostViews() {
		if strings.Contains(strings.ToLower(item.Name), keyword) ||
			strings.Contains(strings.ToLower(item.HostName), keyword) ||
			strings.Contains(strings.ToLower(item.Platform), keyword) ||
			strings.Contains(strings.ToLower(item.Project), keyword) ||
			strings.Contains(strings.ToLower(strings.Join(item.DataAliases, ",")), keyword) ||
			strings.Contains(strings.ToLower(item.GrpcKey), keyword) {
			result = append(result, item)
		}
	}
	return result, nil
}

func (s *GrpcClientProxyService) GetClientHost(ctx context.Context, hostID string) (GrpcClientHostView, error) {
	_ = ctx
	session, err := s.findSession(hostID)
	if err != nil {
		return GrpcClientHostView{}, err
	}
	return sessionToHostView(session), nil
}

func (s *GrpcClientProxyService) RemoveClient(hostID string) bool {
	session, err := s.findSession(hostID)
	if err != nil || session == nil {
		return false
	}
	s.manager.Remove(session.Key)
	return true
}

func (s *GrpcClientProxyService) GetClientDataDirs(ctx context.Context, hostID string) ([]GrpcClientDataDirView, error) {
	session, err := s.findSession(hostID)
	if err != nil {
		return nil, err
	}

	out := make([]GrpcClientDataDirView, 0, len(session.DataSources))
	for _, source := range session.DataSources {
		if !source.Enabled || strings.TrimSpace(source.Alias) == "" {
			continue
		}
		out = append(out, GrpcClientDataDirView{
			Alias:    source.Alias,
			DataDir:   source.Alias,
			HostName: session.HostName,
			HostURL:  session.Key,
			Enabled:  true,
		})
	}
	return out, nil
}

func (s *GrpcClientProxyService) ListDataFiles(ctx context.Context, hostID string, dir string) ([]GrpcFileInfoView, error) {
	session, err := s.findSession(hostID)
	if err != nil {
		return nil, err
	}

	if strings.TrimSpace(dir) == "" {
		out := make([]GrpcFileInfoView, 0, len(session.DataSources))
		index := 0
		for _, source := range session.DataSources {
			if !source.Enabled || strings.TrimSpace(source.Alias) == "" {
				continue
			}
			out = append(out, GrpcFileInfoView{
				Name:    source.Alias,
				Path:    fmt.Sprintf("root://%d", index),
				ModTime: time.Now(),
				IsDir:   true,
				IsData:   false,
				Alias:   source.Alias,
			})
			index++
		}
		return out, nil
	}

	alias, rel := parseAliasAndRelativePath(dir, session)
	response, err := s.sendQuery(ctx, session, &agentv1.QueryRequest{
		Alias:        alias,
		RelativePath: rel,
		Mode:         "list",
	})
	if err != nil {
		return nil, err
	}

	return decodeFileListResponse(response)
}

func (s *GrpcClientProxyService) ReadDataFile(ctx context.Context, req ReadDataRequest) (*ReadDataResult, error) {
	session, err := s.findSession(req.HostID)
	if err != nil {
		return nil, err
	}

	alias, rel := parseAliasAndRelativePath(req.FilePath, session)
	response, err := s.sendQuery(ctx, session, &agentv1.QueryRequest{
		Alias:        alias,
		RelativePath: rel,
		Lines:        int32(req.Lines),
		FromEnd:      req.FromEnd,
		Cursor:       req.Cursor,
		Mode:         req.Mode,
	})
	if err != nil {
		return nil, err
	}

	return &ReadDataResult{
		Lines:      append([]string(nil), response.GetLines()...),
		TotalLines: -1,
		HasMore:    response.GetHasMore(),
		StartLine:  0,
		EndLine:    0,
		Cursor:     response.GetCursor(),
		FileSize:   response.GetFileSize(),
	}, nil
}

func (s *GrpcClientProxyService) DownloadDataFile(ctx context.Context, req DownloadDataRequest) error {
	if req.Writer == nil {
		return errors.New("download writer is nil")
	}

	session, err := s.findSession(req.HostID)
	if err != nil {
		return err
	}

	alias, rel := parseAliasAndRelativePath(req.FilePath, session)
	requestID := newGRPCRequestID()
	chunkCh := make(chan grpcDownloadChunk, downloadChunkChannelSize)

	pendingDownloadMu.Lock()
	pendingDownloads[requestID] = chunkCh
	pendingDownloadMu.Unlock()
	defer func() {
		pendingDownloadMu.Lock()
		delete(pendingDownloads, requestID)
		pendingDownloadMu.Unlock()
		close(chunkCh)
	}()

	sendCtx, sendCancel := context.WithTimeout(ctx, defaultGrpcSendTimeout)
	err = session.SendToClientContext(sendCtx, &agentv1.ServerMessage{
		RequestId: requestID,
		Payload: &agentv1.ServerMessage_QueryRequest{
			QueryRequest: &agentv1.QueryRequest{
				Alias:        alias,
				RelativePath: rel,
				Mode:         "download",
			},
		},
	})
	sendCancel()
	if err != nil {
		return err
	}

	requestCtx := ctx
	if _, ok := ctx.Deadline(); !ok {
		var cancel context.CancelFunc
		requestCtx, cancel = context.WithTimeout(ctx, 30*time.Minute)
		defer cancel()
	}

	for {
		select {
		case <-requestCtx.Done():
			return requestCtx.Err()

		case chunk := <-chunkCh:
			if chunk.err != nil {
				return chunk.err
			}
			if len(chunk.data) > 0 {
				if _, err := req.Writer.Write(chunk.data); err != nil {
					return err
				}
			}
			if chunk.eof {
				return nil
			}
		}
	}
}

func (s *GrpcClientProxyService) allHostViews() []GrpcClientHostView {
	sessions := s.manager.List()
	out := make([]GrpcClientHostView, 0, len(sessions))
	for _, session := range sessions {
		out = append(out, sessionToHostView(session))
	}

	sort.Slice(out, func(i, j int) bool {
		if out[i].Platform == out[j].Platform {
			if out[i].Project == out[j].Project {
				leftAlias := firstString(out[i].DataAliases)
				rightAlias := firstString(out[j].DataAliases)
				if leftAlias == rightAlias {
					return out[i].HostName < out[j].HostName
				}
				return leftAlias < rightAlias
			}
			return out[i].Project < out[j].Project
		}
		return out[i].Platform < out[j].Platform
	})

	return out
}

func (s *GrpcClientProxyService) findSession(hostID string) (*GrpcClientSession, error) {
	hostID = strings.TrimSpace(hostID)
	if hostID == "" {
		return nil, errors.New("hostId is empty")
	}

	session, ok := s.manager.FindLive(hostID, func(item *GrpcClientSession) []string {
		view := sessionToHostView(item)
		return []string{
			view.ClientID,
			fmt.Sprint(view.ID),
			view.GrpcKey,
			view.HostName,
		}
	})
	if !ok || session == nil {
		return nil, fmt.Errorf("gRPC 客户端不存在或已离线: %s", hostID)
	}

	return session, nil
}

func sessionToHostView(session *GrpcClientSession) GrpcClientHostView {
	if session == nil {
		return GrpcClientHostView{}
	}

	dataAliases := make([]string, 0, len(session.DataSources))
	for _, source := range session.DataSources {
		if source.Enabled && strings.TrimSpace(source.Alias) != "" {
			dataAliases = append(dataAliases, source.Alias)
		}
	}
	sort.Strings(dataAliases)

	name := session.HostName
	clientID := strings.Join([]string{session.Platform, session.HostName}, "-")
	url := "grpc://" + session.Key

	return GrpcClientHostView{
		ID:            stableUintID(session.Key),
		Name:          name,
		Url:           url,
		Comments:      fmt.Sprintf("platform=%s project=%s dataAliases=%s hostName=%s status=%s grpc=%s", session.Platform, session.Project, strings.Join(dataAliases, ","), session.HostName, session.Status, session.Key),
		GrpcKey:       session.Key,
		ClientID:      clientID,
		Platform:      session.Platform,
		Project:       session.Project,
		HostName:      session.HostName,
		ClientVersion: session.ClientVersion,
		Users:         append([]string(nil), session.Users...),
		ExampleProjectUsers:   append([]string(nil), session.Users...),
		DataAliases:    dataAliases,
		DataDirs:       append([]GrpcDataSource(nil), session.DataSources...),
		Status:        session.Status,
		RegisteredAt:  formatTimeForView(session.RegisteredAt),
		LastSeenAt:    formatTimeForView(session.LastSeenAt),
		HeartbeatAt:   formatTimeForView(session.LastSeenAt),
	}
}

func parseAliasAndRelativePath(path string, session *GrpcClientSession) (string, string) {
	path = strings.TrimSpace(path)
	if path == "" {
		return firstEnabledAlias(session), ""
	}

	if strings.HasPrefix(path, "root://") {
		indexText := strings.TrimPrefix(path, "root://")
		rel := ""
		if slash := strings.Index(indexText, "/"); slash >= 0 {
			rel = strings.TrimPrefix(indexText[slash+1:], "/")
			indexText = indexText[:slash]
		}

		var index int
		_, _ = fmt.Sscanf(indexText, "%d", &index)
		alias := aliasByIndex(session, index)
		return alias, rel
	}

	path = strings.TrimPrefix(path, "/")
	alias := path
	rel := ""
	if slash := strings.Index(path, "/"); slash >= 0 {
		alias = path[:slash]
		rel = path[slash+1:]
	}

	if strings.TrimSpace(alias) == "" {
		alias = firstEnabledAlias(session)
	}

	return alias, rel
}

func firstString(values []string) string {
	if len(values) == 0 {
		return ""
	}
	return strings.TrimSpace(values[0])
}

func firstEnabledAlias(session *GrpcClientSession) string {
	if session == nil {
		return ""
	}
	for _, source := range session.DataSources {
		if source.Enabled && strings.TrimSpace(source.Alias) != "" {
			return source.Alias
		}
	}
	return ""
}

func aliasByIndex(session *GrpcClientSession, index int) string {
	if session == nil || index < 0 {
		return ""
	}
	current := 0
	for _, source := range session.DataSources {
		if !source.Enabled || strings.TrimSpace(source.Alias) == "" {
			continue
		}
		if current == index {
			return source.Alias
		}
		current++
	}
	return firstEnabledAlias(session)
}

func (s *GrpcClientProxyService) sendQuery(ctx context.Context, session *GrpcClientSession, req *agentv1.QueryRequest) (*agentv1.QueryResponse, error) {
	requestID := newGRPCRequestID()
	resultCh := make(chan grpcQueryResult, 1)

	pendingQueryMu.Lock()
	pendingQueries[requestID] = resultCh
	pendingQueryMu.Unlock()
	defer func() {
		pendingQueryMu.Lock()
		delete(pendingQueries, requestID)
		pendingQueryMu.Unlock()
	}()

	requestCtx := ctx
	if _, ok := ctx.Deadline(); !ok {
		var cancel context.CancelFunc
		requestCtx, cancel = context.WithTimeout(ctx, defaultGrpcRequestTimeout)
		defer cancel()
	}

	sendCtx, sendCancel := context.WithTimeout(requestCtx, defaultGrpcSendTimeout)
	err := session.SendToClientContext(sendCtx, &agentv1.ServerMessage{
		RequestId: requestID,
		Payload: &agentv1.ServerMessage_QueryRequest{
			QueryRequest: req,
		},
	})
	sendCancel()
	if err != nil {
		return nil, err
	}

	select {
	case <-requestCtx.Done():
		return nil, requestCtx.Err()

	case result := <-resultCh:
		if result.err != nil {
			return nil, result.err
		}
		if result.response == nil {
			return nil, errors.New("empty grpc data query response")
		}
		return result.response, nil
	}
}

func decodeFileListResponse(response *agentv1.QueryResponse) ([]GrpcFileInfoView, error) {
	if response == nil || len(response.GetLines()) == 0 {
		return []GrpcFileInfoView{}, nil
	}

	var list []GrpcFileInfoView
	if err := json.Unmarshal([]byte(response.GetLines()[0]), &list); err != nil {
		return nil, fmt.Errorf("decode grpc file list response: %w", err)
	}
	return list, nil
}

func newGRPCRequestID() string {
	buffer := make([]byte, 16)
	if _, err := rand.Read(buffer); err != nil {
		return fmt.Sprintf("req-%d", time.Now().UnixNano())
	}
	return "req-" + hex.EncodeToString(buffer)
}

func stableUintID(value string) uint {
	h := fnv.New32a()
	_, _ = h.Write([]byte(value))
	return uint(h.Sum32())
}

func formatTimeForView(t time.Time) string {
	if t.IsZero() {
		return ""
	}
	return t.UTC().Format(time.RFC3339)
}

7.5 grpc_session_event_service.go

package exampleProject

import (
	"fmt"
	"sync"
	"time"
)

const (
	GrpcSessionEventTypeProjectTreeChanged = "project_tree_changed"

	GrpcSessionEventClientOnline  = "client_online"
	GrpcSessionEventClientOffline = "client_offline"
	GrpcSessionEventClientTimeout = "client_timeout"
	GrpcSessionEventClientReplace = "client_replace"
	GrpcSessionEventClientUpdate  = "client_update"
)

type GrpcSessionEvent struct {
	Type string `json:"type"`

	Event string `json:"event"`

	SessionKey string `json:"sessionKey"`

	Platform string `json:"platform"`

	// Project 是子项目名称,来自客户端 example-project.project,用于前端项目树渲染。
	// 数据查询 alias 请使用 DataAliases。
	Project string `json:"project"`

	DataAliases []string `json:"dataAliases"`

	HostName string `json:"hostName"`

	Status string `json:"status"`

	Time string `json:"time"`
}

type GrpcSessionEventSubscription struct {
	id uint64

	C <-chan GrpcSessionEvent
}

type grpcSessionEventSubscriber struct {
	id uint64

	ch chan GrpcSessionEvent
}

type GrpcSessionEventHub struct {
	mu sync.RWMutex

	nextID uint64

	subscribers map[uint64]*grpcSessionEventSubscriber
}

var defaultGrpcSessionEventHub = NewGrpcSessionEventHub()

func GetGrpcSessionEventHub() *GrpcSessionEventHub {
	return defaultGrpcSessionEventHub
}

func NewGrpcSessionEventHub() *GrpcSessionEventHub {
	return &GrpcSessionEventHub{
		subscribers: make(map[uint64]*grpcSessionEventSubscriber),
	}
}

func (h *GrpcSessionEventHub) Subscribe(buffer int) *GrpcSessionEventSubscription {
	if h == nil {
		return nil
	}

	if buffer <= 0 {
		buffer = 16
	}

	h.mu.Lock()
	defer h.mu.Unlock()

	h.nextID++
	id := h.nextID

	subscriber := &grpcSessionEventSubscriber{
		id: id,
		ch: make(chan GrpcSessionEvent, buffer),
	}

	h.subscribers[id] = subscriber

	return &GrpcSessionEventSubscription{
		id: id,
		C:  subscriber.ch,
	}
}

func (h *GrpcSessionEventHub) Unsubscribe(subscription *GrpcSessionEventSubscription) {
	if h == nil || subscription == nil {
		return
	}

	h.mu.Lock()
	defer h.mu.Unlock()

	subscriber, exists := h.subscribers[subscription.id]
	if !exists || subscriber == nil {
		return
	}

	delete(h.subscribers, subscription.id)
	close(subscriber.ch)
}

// SubscriberCount 返回当前订阅者数量,主要用于排查 WebSocket 是否真正订阅到了事件。
func (h *GrpcSessionEventHub) SubscriberCount() in

圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言