Documentation
¶
Index ¶
- Constants
- Variables
- func AuthInterceptor(apiToken, jwtSecret string) grpc.UnaryServerInterceptor
- func GetNodeIDFromContext(ctx context.Context) uint32
- func GetPeerAddr(ctx context.Context) string
- func GetUserAdminFromContext(ctx context.Context) (bool, bool)
- func GetUserEmailFromContext(ctx context.Context) (string, bool)
- func GetUserIDFromContext(ctx context.Context) (uint, bool)
- func LoggingInterceptor() grpc.UnaryServerInterceptor
- func RegisterNodeLogServiceServer(s grpc.ServiceRegistrar, srv NodeLogServiceServer)
- func Run(cfg *ServerConfig) error
- func SetNodeIDToContext(ctx context.Context, nodeID uint32) context.Context
- func SetUserAdminToContext(ctx context.Context, isAdmin bool) context.Context
- func SetUserEmailToContext(ctx context.Context, email string) context.Context
- func SetUserIDToContext(ctx context.Context, userID uint) context.Context
- func StreamAuthInterceptor(apiToken, jwtSecret string) grpc.StreamServerInterceptor
- func StreamLoggingInterceptor() grpc.StreamServerInterceptor
- type AgentControlConnection
- type AgentControlGRPCServer
- type AgentControlManager
- func (m *AgentControlManager) AddObservedStateHandler(handler ObservedStateHandler)
- func (m *AgentControlManager) CancelOperation(ctx context.Context, nodeID uint32, operationID string, revision uint64) error
- func (m *AgentControlManager) Connection(nodeID uint32) (AgentControlSnapshot, bool)
- func (m *AgentControlManager) DesiredRevision(nodeID uint32) uint64
- func (m *AgentControlManager) DispatchOperation(ctx context.Context, nodeID uint32, operation *agentv1pb.DesiredOperation) (*agentv1pb.OperationAck, error)
- func (m *AgentControlManager) ObservedState(nodeID uint32) (*agentv1pb.ObservedState, bool)
- type AgentControlSnapshot
- type ConfigSyncGRPCServer
- func (s *ConfigSyncGRPCServer) ConfigChanges(stream pb.ConfigSyncService_ConfigChangesServer) error
- func (s *ConfigSyncGRPCServer) FullSync(ctx context.Context, req *pb.ConfigSyncRequest) (*pb.ConfigSyncResponse, error)
- func (s *ConfigSyncGRPCServer) SyncConfig(ctx context.Context, req *pb.ConfigSyncRequest) (*pb.ConfigSyncResponse, error)
- type DesiredOperationDispatcher
- type HealthGRPCServer
- type KernelOperationBridge
- type NodeConnection
- type NodeConnectionManager
- func (m *NodeConnectionManager) GetActiveNodes() []uint32
- func (m *NodeConnectionManager) GetConfigVersion(nodeID uint32) int64
- func (m *NodeConnectionManager) GetConnection(nodeID uint32) (*NodeConnection, bool)
- func (m *NodeConnectionManager) GetNodeConfigVersion(nodeID uint32) int64
- func (m *NodeConnectionManager) IsConfigChanged(nodeID uint32, currentVer int64) bool
- func (m *NodeConnectionManager) NotifyConfigChange(nodeID uint32)
- func (m *NodeConnectionManager) Register(nodeID uint32, addr string)
- func (m *NodeConnectionManager) SetNodeConfigVersion(nodeID uint32, version int64)
- func (m *NodeConnectionManager) Unregister(nodeID uint32)
- func (m *NodeConnectionManager) UpdateConfigVersion(nodeID uint32, version int64)
- func (m *NodeConnectionManager) UpdateLastSeen(nodeID uint32)
- type NodeGRPCServer
- func (s *NodeGRPCServer) GetConfig(ctx context.Context, req *pb.NodeConfigRequest) (*pb.NodeConfigResponse, error)
- func (s *NodeGRPCServer) Register(ctx context.Context, req *pb.NodeRegisterRequest) (*pb.NodeRegisterResponse, error)
- func (s *NodeGRPCServer) ReportStatus(ctx context.Context, req *pb.NodeStatusRequest) (*pb.StatusResponse, error)
- func (s *NodeGRPCServer) StatusStream(stream pb.NodeService_StatusStreamServer) error
- type NodeLogGRPCServer
- type NodeLogServiceClient
- type NodeLogServiceServer
- type ObservedStateHandler
- type Server
- type ServerConfig
- type TrafficGRPCServer
- func (s *TrafficGRPCServer) OnlineStream(stream pb.TrafficService_OnlineStreamServer) error
- func (s *TrafficGRPCServer) ReportOnline(ctx context.Context, req *pb.OnlineReportRequest) (*pb.StatusResponse, error)
- func (s *TrafficGRPCServer) ReportTraffic(ctx context.Context, req *pb.TrafficReportRequest) (*pb.TrafficReportResponse, error)
- func (s *TrafficGRPCServer) TrafficStream(stream pb.TrafficService_TrafficStreamServer) error
- type UserGRPCServer
Constants ¶
const (
AgentProtocolVersion = agentcontrol.ProtocolV1
)
const NodeLogService_ReportLogs_FullMethodName = "/v2board.NodeLogService/ReportLogs"
Variables ¶
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 ¶
GetNodeIDFromContext 从上下文获取节点ID(需要先在拦截器中设置)
func GetUserAdminFromContext ¶
GetUserAdminFromContext 从上下文获取用户管理员状态
func GetUserEmailFromContext ¶
GetUserEmailFromContext 从上下文获取用户邮箱
func GetUserIDFromContext ¶
GetUserIDFromContext 从上下文获取用户ID
func LoggingInterceptor ¶
func LoggingInterceptor() grpc.UnaryServerInterceptor
LoggingInterceptor 日志拦截器
func RegisterNodeLogServiceServer ¶
func RegisterNodeLogServiceServer(s grpc.ServiceRegistrar, srv NodeLogServiceServer)
func SetNodeIDToContext ¶
SetNodeIDToContext 设置节点ID到上下文
func SetUserAdminToContext ¶
SetUserAdminToContext 设置用户管理员状态到上下文
func SetUserEmailToContext ¶
SetUserEmailToContext 设置用户邮箱到上下文
func SetUserIDToContext ¶
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 ¶
func (s *AgentControlGRPCServer) ControlStream(stream agentv1pb.AgentControlService_ControlStreamServer) error
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 ¶
func (s *ConfigSyncGRPCServer) ConfigChanges(stream pb.ConfigSyncService_ConfigChangesServer) error
ConfigChanges 双向流:配置变更实时推送
func (*ConfigSyncGRPCServer) FullSync ¶
func (s *ConfigSyncGRPCServer) FullSync(ctx context.Context, req *pb.ConfigSyncRequest) (*pb.ConfigSyncResponse, error)
FullSync 全量同步 - 节点启动时获取完整配置
func (*ConfigSyncGRPCServer) SyncConfig ¶
func (s *ConfigSyncGRPCServer) SyncConfig(ctx context.Context, req *pb.ConfigSyncRequest) (*pb.ConfigSyncResponse, error)
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 ¶
func (s *HealthGRPCServer) Check(ctx context.Context, req *pb.HealthCheckRequest) (*pb.HealthCheckResponse, error)
Check 简单健康检查
func (*HealthGRPCServer) Watch ¶
func (s *HealthGRPCServer) Watch(stream pb.HealthService_WatchServer) error
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)
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 (*NodeGRPCServer) GetConfig ¶
func (s *NodeGRPCServer) GetConfig(ctx context.Context, req *pb.NodeConfigRequest) (*pb.NodeConfigResponse, error)
GetConfig 获取节点配置
func (*NodeGRPCServer) Register ¶
func (s *NodeGRPCServer) Register(ctx context.Context, req *pb.NodeRegisterRequest) (*pb.NodeRegisterResponse, error)
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 ¶
func (s *NodeLogGRPCServer) ReportLogs(ctx context.Context, req *dynamicpb.Message) (*pb.StatusResponse, error)
type NodeLogServiceClient ¶
type NodeLogServiceClient interface {
ReportLogs(ctx context.Context, in *dynamicpb.Message, opts ...grpc.CallOption) (*pb.StatusResponse, error)
}
func NewNodeLogServiceClient ¶
func NewNodeLogServiceClient(cc grpc.ClientConnInterface) NodeLogServiceClient
type NodeLogServiceServer ¶
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 (*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 ¶
GracefulShutdown 优雅关闭
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 服务器配置
type TrafficGRPCServer ¶
type TrafficGRPCServer struct {
pb.UnimplementedTrafficServiceServer
// contains filtered or unexported fields
}
TrafficGRPCServer 流量 gRPC 服务实现
func NewTrafficGRPCServer ¶
func NewTrafficGRPCServer() *TrafficGRPCServer
NewTrafficGRPCServer 创建流量 gRPC 服务
func (*TrafficGRPCServer) OnlineStream ¶
func (s *TrafficGRPCServer) OnlineStream(stream pb.TrafficService_OnlineStreamServer) error
OnlineStream 双向流:实时在线状态
func (*TrafficGRPCServer) ReportOnline ¶
func (s *TrafficGRPCServer) ReportOnline(ctx context.Context, req *pb.OnlineReportRequest) (*pb.StatusResponse, error)
ReportOnline 上报在线状态
func (*TrafficGRPCServer) ReportTraffic ¶
func (s *TrafficGRPCServer) ReportTraffic(ctx context.Context, req *pb.TrafficReportRequest) (*pb.TrafficReportResponse, error)
ReportTraffic 批量上报流量
func (*TrafficGRPCServer) TrafficStream ¶
func (s *TrafficGRPCServer) TrafficStream(stream pb.TrafficService_TrafficStreamServer) error
TrafficStream 双向流:实时流量上报
type UserGRPCServer ¶
type UserGRPCServer struct {
pb.UnimplementedUserServiceServer
// contains filtered or unexported fields
}
UserGRPCServer 用户 gRPC 服务实现
func (*UserGRPCServer) GetUsers ¶
func (s *UserGRPCServer) GetUsers(ctx context.Context, req *pb.UserListRequest) (*pb.UserListResponse, error)
GetUsers 获取用户列表
func (*UserGRPCServer) UserChanges ¶
func (s *UserGRPCServer) UserChanges(stream pb.UserService_UserChangesServer) error
UserChanges 双向流:用户变更实时推送