grpc

package
v4.0.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 20, 2026 License: MIT Imports: 41 Imported by: 0

Documentation

Index

Constants

View Source
const (
	AgentProtocolVersion = agentcontrol.ProtocolV1
)
View Source
const NodeLogService_ReportLogs_FullMethodName = "/v2board.NodeLogService/ReportLogs"

Variables

View Source
var NodeLogService_ServiceDesc = grpc.ServiceDesc{
	ServiceName: "v2board.NodeLogService",
	HandlerType: (*NodeLogServiceServer)(nil),
	Methods: []grpc.MethodDesc{
		{
			MethodName: "ReportLogs",
			Handler:    _NodeLogService_ReportLogs_Handler,
		},
	},
	Metadata: "api/grpc/v2board.proto",
}

Functions

func AuthInterceptor

func AuthInterceptor(apiToken, jwtSecret string) grpc.UnaryServerInterceptor

AuthInterceptor 认证拦截器: 先校验节点自身 x-api-key, 否则回退到全局 API Token / JWT 两种认证方式 (管理端/兼容旧调用)

func GetNodeIDFromContext

func GetNodeIDFromContext(ctx context.Context) uint32

GetNodeIDFromContext 从上下文获取节点ID(需要先在拦截器中设置)

func GetPeerAddr

func GetPeerAddr(ctx context.Context) string

GetPeerAddr 从上下文获取客户端地址

func GetUserAdminFromContext

func GetUserAdminFromContext(ctx context.Context) (bool, bool)

GetUserAdminFromContext 从上下文获取用户管理员状态

func GetUserEmailFromContext

func GetUserEmailFromContext(ctx context.Context) (string, bool)

GetUserEmailFromContext 从上下文获取用户邮箱

func GetUserIDFromContext

func GetUserIDFromContext(ctx context.Context) (uint, bool)

GetUserIDFromContext 从上下文获取用户ID

func LoggingInterceptor

func LoggingInterceptor() grpc.UnaryServerInterceptor

LoggingInterceptor 日志拦截器

func RegisterNodeLogServiceServer

func RegisterNodeLogServiceServer(s grpc.ServiceRegistrar, srv NodeLogServiceServer)

func Run

func Run(cfg *ServerConfig) error

Run 运行服务器(阻塞)

func SetNodeIDToContext

func SetNodeIDToContext(ctx context.Context, nodeID uint32) context.Context

SetNodeIDToContext 设置节点ID到上下文

func SetUserAdminToContext

func SetUserAdminToContext(ctx context.Context, isAdmin bool) context.Context

SetUserAdminToContext 设置用户管理员状态到上下文

func SetUserEmailToContext

func SetUserEmailToContext(ctx context.Context, email string) context.Context

SetUserEmailToContext 设置用户邮箱到上下文

func SetUserIDToContext

func SetUserIDToContext(ctx context.Context, userID uint) context.Context

SetUserIDToContext 设置用户ID到上下文

func StreamAuthInterceptor

func StreamAuthInterceptor(apiToken, jwtSecret string) grpc.StreamServerInterceptor

StreamAuthInterceptor 流式认证拦截器: 先校验节点自身 x-api-key, 否则回退到 全局 API Token / JWT 两种认证方式 (管理端/兼容旧调用)

func StreamLoggingInterceptor

func StreamLoggingInterceptor() grpc.StreamServerInterceptor

StreamLoggingInterceptor 流式日志拦截器

Types

type AgentControlConnection

type AgentControlConnection struct {
	NodeID       uint32
	SessionID    string
	AgentVersion string
	InstanceID   string
	Capabilities []*agentv1pb.Capability
	ConnectedAt  time.Time
	LastSeen     time.Time
	DesiredRev   uint64
	ObservedRev  uint64
	// contains filtered or unexported fields
}

AgentControlConnection is the current bidirectional control stream for a node.

type AgentControlGRPCServer

type AgentControlGRPCServer struct {
	agentv1pb.UnimplementedAgentControlServiceServer
	// contains filtered or unexported fields
}

AgentControlGRPCServer implements the Agent-first control stream.

func NewAgentControlGRPCServer

func NewAgentControlGRPCServer(manager *AgentControlManager) *AgentControlGRPCServer

func (*AgentControlGRPCServer) ControlStream

type AgentControlManager

type AgentControlManager struct {
	// contains filtered or unexported fields
}

AgentControlManager owns live streams and correlates desired operations with ACKs.

func GetAgentControlManager

func GetAgentControlManager() *AgentControlManager

func NewAgentControlManager

func NewAgentControlManager() *AgentControlManager

func (*AgentControlManager) AddObservedStateHandler

func (m *AgentControlManager) AddObservedStateHandler(handler ObservedStateHandler)

func (*AgentControlManager) CancelOperation

func (m *AgentControlManager) CancelOperation(ctx context.Context, nodeID uint32, operationID string, revision uint64) error

CancelOperation sends an out-of-band cancellation for an operation already present in desired state. It intentionally does not allocate another revision or desired record; the Agent reports the original operation's terminal observation after cancelling its execution context.

func (*AgentControlManager) Connection

func (m *AgentControlManager) Connection(nodeID uint32) (AgentControlSnapshot, bool)

func (*AgentControlManager) DesiredRevision

func (m *AgentControlManager) DesiredRevision(nodeID uint32) uint64

func (*AgentControlManager) DispatchOperation

func (m *AgentControlManager) DispatchOperation(ctx context.Context, nodeID uint32, operation *agentv1pb.DesiredOperation) (*agentv1pb.OperationAck, error)

DispatchOperation pushes a desired operation and waits only for receipt ACK. Runtime completion is reported separately through ObservedState.

func (*AgentControlManager) ObservedState

func (m *AgentControlManager) ObservedState(nodeID uint32) (*agentv1pb.ObservedState, bool)

type AgentControlSnapshot

type AgentControlSnapshot struct {
	NodeID       uint32    `json:"node_id"`
	SessionID    string    `json:"session_id"`
	AgentVersion string    `json:"agent_version"`
	InstanceID   string    `json:"instance_id"`
	Capabilities []string  `json:"capabilities"`
	ConnectedAt  time.Time `json:"connected_at"`
	LastSeen     time.Time `json:"last_seen"`
	DesiredRev   uint64    `json:"desired_revision"`
	ObservedRev  uint64    `json:"observed_revision"`
}

AgentControlSnapshot is a read-only view used by control-plane services.

type ConfigSyncGRPCServer

type ConfigSyncGRPCServer struct {
	pb.UnimplementedConfigSyncServiceServer
	// contains filtered or unexported fields
}

ConfigSyncGRPCServer 配置同步 gRPC 服务实现

func NewConfigSyncGRPCServer

func NewConfigSyncGRPCServer() *ConfigSyncGRPCServer

NewConfigSyncGRPCServer 创建配置同步 gRPC 服务

func (*ConfigSyncGRPCServer) ConfigChanges

ConfigChanges 双向流:配置变更实时推送

func (*ConfigSyncGRPCServer) FullSync

FullSync 全量同步 - 节点启动时获取完整配置

func (*ConfigSyncGRPCServer) SyncConfig

SyncConfig 请求配置同步 - 根据 lastSyncTime 判断是否有变更

type DesiredOperationDispatcher

type DesiredOperationDispatcher interface {
	DispatchOperation(context.Context, uint32, *agentv1pb.DesiredOperation) (*agentv1pb.OperationAck, error)
	CancelOperation(context.Context, uint32, string, uint64) error
}

DesiredOperationDispatcher lets durable runtime workers push desired state without depending on the gRPC stream implementation.

type HealthGRPCServer

type HealthGRPCServer struct {
	pb.UnimplementedHealthServiceServer
}

HealthGRPCServer 健康检查 gRPC 服务实现

func NewHealthGRPCServer

func NewHealthGRPCServer() *HealthGRPCServer

NewHealthGRPCServer 创建健康检查 gRPC 服务

func (*HealthGRPCServer) Check

Check 简单健康检查

func (*HealthGRPCServer) Watch

Watch 双向流:持续健康检查

type KernelOperationBridge

type KernelOperationBridge struct {
	// contains filtered or unexported fields
}

KernelOperationBridge joins durable Kernel operations to the ephemeral Agent stream. The database stays authoritative: the stream supplies only the active session and delivery transport.

func NewKernelOperationBridge

func NewKernelOperationBridge(db *gorm.DB, stream kernelOperationStream) (*KernelOperationBridge, error)

func (*KernelOperationBridge) RunOnce

func (b *KernelOperationBridge) RunOnce(ctx context.Context) (int, error)

RunOnce dispatches all pending node operations and recovers dispatch attempts whose Control process stopped before an ACK could be persisted.

func (*KernelOperationBridge) Start

func (b *KernelOperationBridge) Start(ctx context.Context, interval time.Duration, reportError func(error))

Start runs durable dispatch until ctx is cancelled. Callers supply logging so this package remains usable by tests and non-server control processes.

type NodeConnection

type NodeConnection struct {
	NodeID     uint32
	LastSeen   time.Time
	RemoteAddr string
	Connection any // 可以是具体的流对象
	// contains filtered or unexported fields
}

NodeConnection 节点连接信息

type NodeConnectionManager

type NodeConnectionManager struct {
	// contains filtered or unexported fields
}

NodeConnectionManager 节点连接管理器

func GetConnectionManager

func GetConnectionManager() *NodeConnectionManager

GetConnectionManager 获取连接管理器

func NewNodeConnectionManager

func NewNodeConnectionManager() *NodeConnectionManager

NewNodeConnectionManager 创建连接管理器

func (*NodeConnectionManager) GetActiveNodes

func (m *NodeConnectionManager) GetActiveNodes() []uint32

GetActiveNodes 获取所有活跃节点

func (*NodeConnectionManager) GetConfigVersion

func (m *NodeConnectionManager) GetConfigVersion(nodeID uint32) int64

GetConfigVersion 获取配置版本

func (*NodeConnectionManager) GetConnection

func (m *NodeConnectionManager) GetConnection(nodeID uint32) (*NodeConnection, bool)

GetConnection 获取节点连接

func (*NodeConnectionManager) GetNodeConfigVersion

func (m *NodeConnectionManager) GetNodeConfigVersion(nodeID uint32) int64

GetNodeConfigVersion 获取节点已推送的配置版本

func (*NodeConnectionManager) IsConfigChanged

func (m *NodeConnectionManager) IsConfigChanged(nodeID uint32, currentVer int64) bool

IsConfigChanged 检查配置是否变更

func (*NodeConnectionManager) NotifyConfigChange

func (m *NodeConnectionManager) NotifyConfigChange(nodeID uint32)

NotifyConfigChange 通知节点配置变更

func (*NodeConnectionManager) Register

func (m *NodeConnectionManager) Register(nodeID uint32, addr string)

Register 注册节点连接

func (*NodeConnectionManager) SetNodeConfigVersion

func (m *NodeConnectionManager) SetNodeConfigVersion(nodeID uint32, version int64)

SetNodeConfigVersion 记录节点已推送的配置版本

func (*NodeConnectionManager) Unregister

func (m *NodeConnectionManager) Unregister(nodeID uint32)

Unregister 注销节点连接

func (*NodeConnectionManager) UpdateConfigVersion

func (m *NodeConnectionManager) UpdateConfigVersion(nodeID uint32, version int64)

UpdateConfigVersion 更新配置版本

func (*NodeConnectionManager) UpdateLastSeen

func (m *NodeConnectionManager) UpdateLastSeen(nodeID uint32)

UpdateLastSeen 更新最后活跃时间

type NodeGRPCServer

type NodeGRPCServer struct {
	pb.UnimplementedNodeServiceServer
	// contains filtered or unexported fields
}

NodeGRPCServer 节点 gRPC 服务实现

func NewNodeGRPCServer

func NewNodeGRPCServer() *NodeGRPCServer

NewNodeGRPCServer 创建节点 gRPC 服务

func (*NodeGRPCServer) GetConfig

GetConfig 获取节点配置

func (*NodeGRPCServer) Register

Register 节点注册

func (*NodeGRPCServer) ReportStatus

func (s *NodeGRPCServer) ReportStatus(ctx context.Context, req *pb.NodeStatusRequest) (*pb.StatusResponse, error)

ReportStatus 上报节点状态

func (*NodeGRPCServer) StatusStream

func (s *NodeGRPCServer) StatusStream(stream pb.NodeService_StatusStreamServer) error

StatusStream 双向流:节点状态实时通信

type NodeLogGRPCServer

type NodeLogGRPCServer struct {
	// contains filtered or unexported fields
}

func NewNodeLogGRPCServer

func NewNodeLogGRPCServer() *NodeLogGRPCServer

func (*NodeLogGRPCServer) ReportLogs

type NodeLogServiceClient

type NodeLogServiceClient interface {
	ReportLogs(ctx context.Context, in *dynamicpb.Message, opts ...grpc.CallOption) (*pb.StatusResponse, error)
}

type NodeLogServiceServer

type NodeLogServiceServer interface {
	ReportLogs(context.Context, *dynamicpb.Message) (*pb.StatusResponse, error)
}

type ObservedStateHandler

type ObservedStateHandler func(nodeID uint32, observed *agentv1pb.ObservedState)

ObservedStateHandler adapts terminal Agent state back into durable job state.

type Server

type Server struct {
	// contains filtered or unexported fields
}

Server gRPC 服务器

func NewServer

func NewServer(cfg *ServerConfig) *Server

NewServer 创建 gRPC 服务器

func (*Server) GetAgentControlManager

func (s *Server) GetAgentControlManager() *AgentControlManager

GetAgentControlManager returns the Agent-first desired/observed connection manager.

func (*Server) GetConnectionManager

func (s *Server) GetConnectionManager() *NodeConnectionManager

GetConnectionManager 获取连接管理器

func (*Server) GracefulShutdown

func (s *Server) GracefulShutdown(ctx context.Context) error

GracefulShutdown 优雅关闭

func (*Server) Start

func (s *Server) Start() error

Start 启动服务器

func (*Server) Stop

func (s *Server) Stop()

Stop 停止服务器

type ServerConfig

type ServerConfig struct {
	// 监听地址
	Host string
	// 监听端口
	Port int
	// API Token (可选,用于认证)
	APIToken string
	// JWT Secret (可选,用于 JWT 认证)
	JWTSecret string
	// TLS 证书/私钥路径 (可选,都为空时使用明文)
	TLSCertFile string
	TLSKeyFile  string
	// Keepalive 时间
	KeepaliveTime time.Duration
	// Keepalive 超时
	KeepaliveTimeout time.Duration
	// 最大连接空闲时间
	MaxConnectionIdle time.Duration
	// 最大连接年龄
	MaxConnectionAge time.Duration
}

ServerConfig gRPC 服务器配置

func DefaultServerConfig

func DefaultServerConfig() *ServerConfig

DefaultServerConfig 默认配置

type TrafficGRPCServer

type TrafficGRPCServer struct {
	pb.UnimplementedTrafficServiceServer
	// contains filtered or unexported fields
}

TrafficGRPCServer 流量 gRPC 服务实现

func NewTrafficGRPCServer

func NewTrafficGRPCServer() *TrafficGRPCServer

NewTrafficGRPCServer 创建流量 gRPC 服务

func (*TrafficGRPCServer) OnlineStream

OnlineStream 双向流:实时在线状态

func (*TrafficGRPCServer) ReportOnline

ReportOnline 上报在线状态

func (*TrafficGRPCServer) ReportTraffic

ReportTraffic 批量上报流量

func (*TrafficGRPCServer) TrafficStream

TrafficStream 双向流:实时流量上报

type UserGRPCServer

type UserGRPCServer struct {
	pb.UnimplementedUserServiceServer
	// contains filtered or unexported fields
}

UserGRPCServer 用户 gRPC 服务实现

func NewUserGRPCServer

func NewUserGRPCServer() *UserGRPCServer

NewUserGRPCServer 创建用户 gRPC 服务

func (*UserGRPCServer) GetUsers

GetUsers 获取用户列表

func (*UserGRPCServer) UserChanges

func (s *UserGRPCServer) UserChanges(stream pb.UserService_UserChangesServer) error

UserChanges 双向流:用户变更实时推送

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL