示例项目:
exampleProject
技术范围:Protocol Buffers、HTTP/2、gRPC 双向流、注册确认、心跳、会话管理、断线重连、请求—响应关联、流式下载。
说明:本文已对原始工程名称、仓库地址和业务标识进行脱敏。代码用于展示设计思路;接入实际项目时,请按目录结构和配置体系调整。
gRPC 是一种基于 HTTP/2 的远程过程调用框架。通信双方先用 Protocol Buffers(protobuf)定义服务、方法和消息结构,再由 protoc 生成客户端桩与服务端接口。业务代码操作的是强类型方法和消息,不需要手工拼装 HTTP 请求或 JSON。
gRPC 的几个核心组成如下:
| 组成 | 作用 |
|---|---|
.proto 文件 |
定义服务、RPC 方法及消息结构,是通信契约 |
protoc 与插件 |
生成 Go 消息类型、客户端桩和服务端接口 |
| HTTP/2 | 提供多路复用、双向流、头部压缩和长连接能力 |
| Channel/ClientConn | 客户端到目标服务的逻辑连接,可承载多个 RPC |
| Stream | 一次流式 RPC;本文的长连接实际是一个持续存在的双向流 |
| Metadata/Status | 传递调用元数据和标准化错误状态 |
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 中持续发送消息,适合注册、心跳、服务端主动下发任务以及客户端异步返回结果。
.proto 到运行时调用.proto 中定义 ExampleProjectControl.Connect 双向流方法和消息类型。protoc 生成 *.pb.go 与 *_grpc.pb.go。ExampleProjectControlServer 接口并注册到 grpc.Server。grpc.ClientConn,通过生成的 Client 调用 Connect。Send、Recv 持续读写,直到 context 取消、网络中断或任一端结束 RPC。需要区分四层状态:
Connect stream 已建立:代表本次双向流 RPC 已开始。RegisterAck(Accepted=true)。因此,Dial 成功不等于业务客户端已经在线。示例只有在收到注册确认后才启动心跳并把状态切换为 registered。
一个 HTTP/2 连接可以复用多个 stream;本文用一个长期存在的 stream 承载控制消息。gRPC Go 允许一个 goroutine 调用 Send、另一个 goroutine 调用 Recv,但不应让多个 goroutine 同时调用同一方向的 Send。因此:
sendMutex 串行化注册、心跳和业务响应;sendLoop 消费 session 发送队列,形成单写者;Connect,开始本次双向流 RPC。RegisterRequest 作为首条消息,携带 platform、project、hostName 等信息。SessionManager.Add(session);存在同 key 旧会话时,幂等关闭旧会话并替换为新会话。RegisterAck(Accepted=true),客户端收到后设置 registered=true。Heartbeat,随后定期发送心跳。QueryRequest 或 ReloadConfig;客户端按请求类型执行处理,查询或下载结果通过 QueryResponse 或 DataChunk 返回。客户端完成 Dial 后调用 Connect 创建双向流,并把 RegisterRequest 作为第一条消息。服务端只接受注册消息作为首包;校验成功后创建 session,以规范化后的 platform:hostName 作为唯一键,再返回注册确认。
注册成功后客户端立即发送一次心跳,随后按服务端下发的间隔定时发送。服务端更新 LastSeenAt;清理器定期扫描,超过超时时间的 session 会被关闭、移除并发布状态事件。
推荐满足:
heartbeat interval << heartbeat timeout
例如心跳间隔 15 秒、超时 45 秒。扫描周期与客户端心跳周期最好使用两个独立配置,不要混用同一个字段。
服务端先用 FindLive 获取真实在线 session,把带有唯一 request_id 的请求写入发送队列。客户端执行示例查询或下载后,携带相同 request_id 返回响应。服务端通过 pendingQueries 或 pendingDownloads 将异步响应分发给原始调用方。
同一客户端重新连接时,新 session 替换旧 session,并关闭旧 session 的 done。旧 Connect handler 退出时会按对象身份执行条件删除,因此不会误删刚建立的新 session。
session.Send 不被关闭,因为多个 goroutine 可能仍在发送;直接关闭会产生 send on closed channel。生命周期结束只关闭幂等的 done 信号,发送方同时监听 done 和 context。
| 当前状态 | 发生的事件 | 下一状态 |
|---|---|---|
| 启动 | 开始连接 | 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,重新连接 |
注册前连续失败采用指数退避,并加入随机抖动,防止大量实例同时重连。一次连接真正注册成功后再断开,下次重连从初始退避开始。
建议保存为 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,不要手工编辑生成文件。
| 文件 | 端 | 核心职责 |
|---|---|---|
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、注册确认、心跳、重连、接收指令并返回结果 |
RegisterAck.Accepted == true 才算注册成功。RegisterRequest。platform:hostName。List 返回不可写快照;下发消息必须通过 FindLive 获取真实 session。Close() 只关闭 done,不关闭可能被并发写入的 Send。request_id 必须全局足够唯一,并在请求完成、超时或取消后清理等待表。以下代码保留原笔记中的完整实现,并统一使用脱敏后的 exampleProject 命名。
grpc_server_service.gopackage 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())
}
grpc_example_project_server_service.gopackage 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)
}
}
grpc_session_manager_service.gopackage 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
}
grpc_client_proxy_service.gopackage 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)
}
grpc_session_event_service.gopackage 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