Documentation
¶
Index ¶
- Constants
- Variables
- func CanTransitionSessionState(fromState SessionState, toState SessionState) bool
- func CanTransitionTunnelState(fromState TunnelState, toState TunnelState) bool
- func ControlFrameTypeForMessageType(messageType pb.ControlMessageType) (uint16, error)
- func ControlMessageTypeForFrameType(frameType uint16) (pb.ControlMessageType, error)
- func DecodeBusinessControlEnvelopeFrame(frame ControlFrame) (pb.ControlEnvelope, error)
- func IsErrorKind(err error, kind ErrorKind) bool
- type BindingInfo
- type BindingType
- type ControlChannel
- type ControlDiagnosticFields
- type ControlFragmentationConfig
- type ControlFrame
- type ControlFrameFragmenter
- type ControlFrameReassembler
- type ControlMessagePriority
- type Error
- type ErrorKind
- type FailedMappingContext
- type HeartbeatMonitor
- func (monitor *HeartbeatMonitor) ObservePeerActivity(now time.Time)
- func (monitor *HeartbeatMonitor) ObserveReceive(now time.Time)
- func (monitor *HeartbeatMonitor) ObserveSend(now time.Time)
- func (monitor *HeartbeatMonitor) ShouldSend(now time.Time) bool
- func (monitor *HeartbeatMonitor) Snapshot(now time.Time) HeartbeatStatus
- func (monitor *HeartbeatMonitor) Start(now time.Time)
- type HeartbeatPolicy
- type HeartbeatStatus
- type IdleTunnelEviction
- type IdleTunnelEvictionReason
- type InMemorySession
- func (session *InMemorySession) BindingInfo() BindingInfo
- func (session *InMemorySession) Close(ctx context.Context, reason error) error
- func (session *InMemorySession) Control() (ControlChannel, error)
- func (session *InMemorySession) Done() <-chan struct{}
- func (session *InMemorySession) Err() error
- func (session *InMemorySession) ID() string
- func (session *InMemorySession) MarkFailed(cause error, retrying bool, giveUp bool) error
- func (session *InMemorySession) Meta() SessionMeta
- func (session *InMemorySession) Open(ctx context.Context) error
- func (session *InMemorySession) ProtocolState() ProtocolState
- func (session *InMemorySession) State() SessionState
- func (session *InMemorySession) TunnelAcceptor() (TunnelAcceptor, error)
- func (session *InMemorySession) TunnelPool() (TunnelPool, error)
- func (session *InMemorySession) TunnelProducer() (TunnelProducer, error)
- type InMemoryTunnelPool
- func (pool *InMemoryTunnelPool) Acquire(ctx context.Context) (Tunnel, error)
- func (pool *InMemoryTunnelPool) ClosedCount() int
- func (pool *InMemoryTunnelPool) Config() TunnelPoolConfig
- func (pool *InMemoryTunnelPool) EvictExpiredIdle(now time.Time) []string
- func (pool *InMemoryTunnelPool) EvictZombieIdle(ctx context.Context, now time.Time) []IdleTunnelEviction
- func (pool *InMemoryTunnelPool) IdleCount() int
- func (pool *InMemoryTunnelPool) InUseCount() int
- func (pool *InMemoryTunnelPool) PutIdle(tunnel Tunnel) error
- func (pool *InMemoryTunnelPool) Recycle(ctx context.Context, tunnel Tunnel) (RecycleResult, error)
- func (pool *InMemoryTunnelPool) RecycledCount() int
- func (pool *InMemoryTunnelPool) Remove(tunnelID string) error
- type KeepalivePolicy
- type PoolDiagnosticFields
- type PoolRefillDiagnosticFields
- type PrioritizedControlChannel
- type PrioritizedControlFrame
- type PriorityControlQueue
- type ProtocolState
- type RecycleResult
- type RefillController
- type RefillControllerConfig
- type RefillRequestDeduplicator
- type RefillResult
- type Session
- type SessionCapabilities
- type SessionDiagnosticFields
- type SessionMeta
- type SessionState
- type TransportMetricsRecorder
- func (recorder *TransportMetricsRecorder) IncBrokenTunnel()
- func (recorder *TransportMetricsRecorder) IncOpenTimeout()
- func (recorder *TransportMetricsRecorder) IncReset()
- func (recorder *TransportMetricsRecorder) ObserveHeartbeatRTT(roundTripTime time.Duration)
- func (recorder *TransportMetricsRecorder) ObservePoolCounts(idleCount int, inUseCount int)
- func (recorder *TransportMetricsRecorder) ObserveRefill(openedCount int, observedAt time.Time)
- func (recorder *TransportMetricsRecorder) Snapshot(now time.Time) TransportMetricsSnapshot
- type TransportMetricsSnapshot
- type Tunnel
- type TunnelAcceptor
- type TunnelDiagnosticFields
- type TunnelHealthProber
- type TunnelMeta
- type TunnelPool
- type TunnelPoolConfig
- type TunnelPoolReport
- type TunnelProducer
- type TunnelRefillReason
- type TunnelRefillRequest
- type TunnelState
Constants ¶
const ( // ControlFrameTypeFragment 是控制面分块帧的保留类型。 ControlFrameTypeFragment uint16 = 0xFFFF // DefaultControlFragmentMaxPayloadSize 是控制面分块的默认单帧上限。 DefaultControlFragmentMaxPayloadSize = 64 * 1024 )
const ( // ControlFrameTypeHeartbeatPing 表示控制面 heartbeat ping 帧。 ControlFrameTypeHeartbeatPing uint16 = 0xFFF0 // ControlFrameTypeHeartbeatPong 表示控制面 heartbeat pong 帧。 ControlFrameTypeHeartbeatPong uint16 = 0xFFF1 )
const ( // ControlFrameTypeConnectorHello 表示 ConnectorHello 业务控制帧。 ControlFrameTypeConnectorHello uint16 = 0x0101 // ControlFrameTypeConnectorWelcome 表示 ConnectorWelcome 业务控制帧。 ControlFrameTypeConnectorWelcome uint16 = 0x0102 // ControlFrameTypeConnectorAuth 表示 ConnectorAuth 业务控制帧。 ControlFrameTypeConnectorAuth uint16 = 0x0103 // ControlFrameTypeConnectorAuthAck 表示 ConnectorAuthAck 业务控制帧。 ControlFrameTypeConnectorAuthAck uint16 = 0x0104 // ControlFrameTypeControlHeartbeat 表示控制面业务心跳帧(区别于 ping/pong 保活帧)。 ControlFrameTypeControlHeartbeat uint16 = 0x0105 // ControlFrameTypePublishService 表示 PublishService 业务控制帧。 ControlFrameTypePublishService uint16 = 0x0110 // ControlFrameTypePublishServiceAck 表示 PublishServiceAck 业务控制帧。 ControlFrameTypePublishServiceAck uint16 = 0x0111 // ControlFrameTypeUnpublishService 表示 UnpublishService 业务控制帧。 ControlFrameTypeUnpublishService uint16 = 0x0112 // ControlFrameTypeUnpublishServiceAck 表示 UnpublishServiceAck 业务控制帧。 ControlFrameTypeUnpublishServiceAck uint16 = 0x0113 // ControlFrameTypeServiceHealthReport 表示 ServiceHealthReport 业务控制帧。 ControlFrameTypeServiceHealthReport uint16 = 0x0114 // ControlFrameTypeTunnelPoolReport 表示 TunnelPoolReport 业务控制帧。 ControlFrameTypeTunnelPoolReport uint16 = 0x0115 // ControlFrameTypeTunnelRefillRequest 表示 TunnelRefillRequest 业务控制帧。 ControlFrameTypeTunnelRefillRequest uint16 = 0x0116 // ControlFrameTypeTunnelDialAnnounce 表示 TunnelDialAnnounce 业务控制帧。 ControlFrameTypeTunnelDialAnnounce uint16 = 0x0117 // ControlFrameTypeRouteAssign 表示 RouteAssign 业务控制帧。 ControlFrameTypeRouteAssign uint16 = 0x0120 // ControlFrameTypeRouteAssignAck 表示 RouteAssignAck 业务控制帧。 ControlFrameTypeRouteAssignAck uint16 = 0x0121 // ControlFrameTypeRouteRevoke 表示 RouteRevoke 业务控制帧。 ControlFrameTypeRouteRevoke uint16 = 0x0122 // ControlFrameTypeRouteRevokeAck 表示 RouteRevokeAck 业务控制帧。 ControlFrameTypeRouteRevokeAck uint16 = 0x0123 // ControlFrameTypeRouteStatusReport 表示 RouteStatusReport 业务控制帧。 ControlFrameTypeRouteStatusReport uint16 = 0x0124 // ControlFrameTypeControlError 表示 ControlError 业务控制帧。 ControlFrameTypeControlError uint16 = 0x01F0 )
Variables ¶
var ( // ErrClosed 表示资源已关闭。 ErrClosed = errors.New("closed") // ErrNotReady 表示资源尚未就绪。 ErrNotReady = errors.New("not ready") // ErrSessionClosed 表示 session 已关闭或不可用。 ErrSessionClosed = errors.New("session closed") // ErrTunnelClosed 表示 tunnel 已关闭。 ErrTunnelClosed = errors.New("tunnel closed") // ErrTunnelBroken 表示 tunnel 已损坏。 ErrTunnelBroken = errors.New("tunnel broken") // ErrTunnelStale 表示 tunnel 探活失败,不应继续分配。 ErrTunnelStale = errors.New("tunnel stale: probe failed") // ErrNoTunnel 表示当前没有可分配 tunnel。 ErrNoTunnel = errors.New("no available tunnel") // ErrTimeout 表示操作超时。 ErrTimeout = errors.New("timeout") // ErrUnsupported 表示能力不支持。 ErrUnsupported = errors.New("unsupported capability") // ErrInvalidArgument 表示参数非法。 ErrInvalidArgument = errors.New("invalid argument") // ErrStateTransition 表示状态迁移非法。 ErrStateTransition = errors.New("invalid state transition") // ErrTunnelNotFound 表示未找到目标 tunnel。 ErrTunnelNotFound = errors.New("tunnel not found") // ErrTunnelInUse 表示 tunnel 已在使用中。 ErrTunnelInUse = errors.New("tunnel in use") // ErrPoolExhausted 表示池容量已达到上限。 ErrPoolExhausted = errors.New("pool exhausted") )
Functions ¶
func CanTransitionSessionState ¶
func CanTransitionSessionState(fromState SessionState, toState SessionState) bool
CanTransitionSessionState 判断 session 状态迁移是否合法。
func CanTransitionTunnelState ¶
func CanTransitionTunnelState(fromState TunnelState, toState TunnelState) bool
CanTransitionTunnelState 判断 tunnel 状态迁移是否合法。
func ControlFrameTypeForMessageType ¶
func ControlFrameTypeForMessageType(messageType pb.ControlMessageType) (uint16, error)
ControlFrameTypeForMessageType 将控制面消息类型映射到控制帧类型。
func ControlMessageTypeForFrameType ¶
func ControlMessageTypeForFrameType(frameType uint16) (pb.ControlMessageType, error)
ControlMessageTypeForFrameType 将业务控制帧类型映射为控制面消息类型。
func DecodeBusinessControlEnvelopeFrame ¶
func DecodeBusinessControlEnvelopeFrame(frame ControlFrame) (pb.ControlEnvelope, error)
DecodeBusinessControlEnvelopeFrame 将业务控制帧解码为控制面封装。
func IsErrorKind ¶
IsErrorKind 判断 error 是否为指定 ErrorKind。
Types ¶
type BindingInfo ¶
type BindingInfo struct {
Type BindingType
Version string
MaxConcurrentStreams int64
KeepalivePolicy KeepalivePolicy
}
BindingInfo 描述当前 session 的 binding 元信息。
func NewBindingInfo ¶
func NewBindingInfo(bindingType BindingType) BindingInfo
NewBindingInfo 创建带默认 keepalive 策略的 BindingInfo。
type BindingType ¶
type BindingType string
BindingType 表示 transport binding 的实现类型。
const ( // BindingTypeGRPCH2 表示 gRPC over HTTP/2 binding。 BindingTypeGRPCH2 BindingType = "grpc_h2" // BindingTypeTCPFramed 表示基于长度前缀帧的 TCP binding。 BindingTypeTCPFramed BindingType = "tcp_framed" // BindingTypeQUICNative 表示原生 QUIC binding。 BindingTypeQUICNative BindingType = "quic_native" // BindingTypeH3Stream 表示 HTTP/3 stream binding。 BindingTypeH3Stream BindingType = "h3_stream" )
func (BindingType) String ¶
func (bindingType BindingType) String() string
String 返回 binding 类型的字符串表示。
type ControlChannel ¶
type ControlChannel interface {
WriteControlFrame(ctx context.Context, frame ControlFrame) error
ReadControlFrame(ctx context.Context) (ControlFrame, error)
Close(ctx context.Context) error
Done() <-chan struct{}
Err() error
}
ControlChannel 定义控制面的最小读写能力。
type ControlDiagnosticFields ¶
type ControlDiagnosticFields struct {
SessionID string
SessionEpoch uint64
Binding string
LastError string
Retryable bool
}
ControlDiagnosticFields 描述控制面错误与重试日志的统一字段。
func BuildControlDiagnosticFields ¶
func BuildControlDiagnosticFields( meta SessionMeta, bindingInfo BindingInfo, lastError error, retryable bool, ) ControlDiagnosticFields
BuildControlDiagnosticFields 从会话元信息构造统一诊断字段。
func BuildControlDiagnosticFieldsForSession ¶
func BuildControlDiagnosticFieldsForSession(session Session, retryable bool) ControlDiagnosticFields
BuildControlDiagnosticFieldsForSession 从 transport.Session 直接提取统一诊断字段。
func (ControlDiagnosticFields) Fields ¶
func (fields ControlDiagnosticFields) Fields() map[string]any
Fields 返回适合结构化日志输出的字段映射。
type ControlFragmentationConfig ¶
type ControlFragmentationConfig struct {
MaxPayloadSize int
ReassemblyTTL time.Duration
MaxInFlightMessage int
}
ControlFragmentationConfig 描述控制面大消息分块与重组配置。
func DefaultControlFragmentationConfig ¶
func DefaultControlFragmentationConfig() ControlFragmentationConfig
DefaultControlFragmentationConfig 返回默认控制面分块配置。
func (ControlFragmentationConfig) NormalizeAndValidate ¶
func (config ControlFragmentationConfig) NormalizeAndValidate() (ControlFragmentationConfig, error)
NormalizeAndValidate 归一化并校验控制面分块配置。
type ControlFrame ¶
ControlFrame 描述控制面的最小帧结构。
func EncodeBusinessControlEnvelopeFrame ¶
func EncodeBusinessControlEnvelopeFrame(envelope pb.ControlEnvelope) (ControlFrame, error)
EncodeBusinessControlEnvelopeFrame 将业务控制面封装编码为控制帧。
type ControlFrameFragmenter ¶
type ControlFrameFragmenter struct {
// contains filtered or unexported fields
}
ControlFrameFragmenter 负责把大控制帧拆分成多个分块帧。
func NewControlFrameFragmenter ¶
func NewControlFrameFragmenter(config ControlFragmentationConfig) (*ControlFrameFragmenter, error)
NewControlFrameFragmenter 创建控制面分块器。
func (*ControlFrameFragmenter) Fragment ¶
func (fragmenter *ControlFrameFragmenter) Fragment(frame ControlFrame) ([]ControlFrame, error)
Fragment 按配置把控制帧拆成一个或多个分块帧。
type ControlFrameReassembler ¶
type ControlFrameReassembler struct {
// contains filtered or unexported fields
}
ControlFrameReassembler 负责把分块控制帧重组成原始控制帧。
func NewControlFrameReassembler ¶
func NewControlFrameReassembler(config ControlFragmentationConfig) (*ControlFrameReassembler, error)
NewControlFrameReassembler 创建控制面分块重组器。
func (*ControlFrameReassembler) Reassemble ¶
func (reassembler *ControlFrameReassembler) Reassemble(frame ControlFrame) (ControlFrame, bool, error)
Reassemble 尝试把单帧或分块帧重组为原始控制帧。
type ControlMessagePriority ¶
type ControlMessagePriority string
ControlMessagePriority 表示控制面消息优先级。
const ( // ControlMessagePriorityHigh 表示高优先级消息,例如 heartbeat、auth、refill。 ControlMessagePriorityHigh ControlMessagePriority = "high" // ControlMessagePriorityNormal 表示常规优先级消息。 ControlMessagePriorityNormal ControlMessagePriority = "normal" // ControlMessagePriorityLow 表示低优先级消息,例如大体积状态同步。 ControlMessagePriorityLow ControlMessagePriority = "low" )
func NormalizeControlMessagePriority ¶
func NormalizeControlMessagePriority(priority ControlMessagePriority) ControlMessagePriority
NormalizeControlMessagePriority 归一化控制面优先级取值。
func RecommendControlFramePriority ¶
func RecommendControlFramePriority(frameType uint16) ControlMessagePriority
RecommendControlFramePriority 按控制帧类型返回建议发送优先级。
type Error ¶
Error 定义 transport 统一错误结构。
type ErrorKind ¶
type ErrorKind string
ErrorKind 表示 transport 错误分类。
const ( // ErrorKindTransport 表示底层传输错误。 ErrorKindTransport ErrorKind = "transport" // ErrorKindProtocol 表示协议顺序或格式错误。 ErrorKindProtocol ErrorKind = "protocol" // ErrorKindReject 表示业务拒绝。 ErrorKindReject ErrorKind = "reject" // ErrorKindTimeout 表示超时类错误。 ErrorKindTimeout ErrorKind = "timeout" // ErrorKindClosed 表示资源关闭类错误。 ErrorKindClosed ErrorKind = "closed" // ErrorKindInternal 表示实现内部错误。 ErrorKindInternal ErrorKind = "internal" )
type FailedMappingContext ¶
type FailedMappingContext struct {
// EverAuthenticated 表示本次 session 是否曾进入 authenticated。
EverAuthenticated bool
// Retrying 表示当前是否处于自动重试/退避重连阶段。
Retrying bool
// GiveUp 表示实现是否已经放弃本次 session 建立流程。
GiveUp bool
}
FailedMappingContext 描述 SessionState=failed 时的映射上下文。
type HeartbeatMonitor ¶
type HeartbeatMonitor struct {
// contains filtered or unexported fields
}
HeartbeatMonitor 负责跟踪控制面 heartbeat 的发送与判死状态。
func NewHeartbeatMonitor ¶
func NewHeartbeatMonitor(policy HeartbeatPolicy) (*HeartbeatMonitor, error)
NewHeartbeatMonitor 创建 heartbeat 状态跟踪器。
func (*HeartbeatMonitor) ObservePeerActivity ¶
func (monitor *HeartbeatMonitor) ObservePeerActivity(now time.Time)
ObservePeerActivity 是 ObserveReceive 的语义化别名,便于 control loop 在收到任意对端控制消息时复用。
func (*HeartbeatMonitor) ObserveReceive ¶
func (monitor *HeartbeatMonitor) ObserveReceive(now time.Time)
ObserveReceive 记录一条来自对端的 heartbeat 或等价存活信号。
func (*HeartbeatMonitor) ObserveSend ¶
func (monitor *HeartbeatMonitor) ObserveSend(now time.Time)
ObserveSend 记录一条 heartbeat 已经真正发送完成。
func (*HeartbeatMonitor) ShouldSend ¶
func (monitor *HeartbeatMonitor) ShouldSend(now time.Time) bool
ShouldSend 判断当前时间点是否已到下一次 heartbeat 发送时机。
func (*HeartbeatMonitor) Snapshot ¶
func (monitor *HeartbeatMonitor) Snapshot(now time.Time) HeartbeatStatus
Snapshot 返回指定时间点的 heartbeat 运行快照。
func (*HeartbeatMonitor) Start ¶
func (monitor *HeartbeatMonitor) Start(now time.Time)
Start 标记控制面 heartbeat 从指定时间点开始生效。
type HeartbeatPolicy ¶
type HeartbeatPolicy struct {
SendInterval time.Duration
MissThreshold int
QueueDelayGrace time.Duration
}
HeartbeatPolicy 描述控制面 heartbeat 的发送与判死策略。
func DefaultHeartbeatPolicy ¶
func DefaultHeartbeatPolicy() HeartbeatPolicy
DefaultHeartbeatPolicy 返回默认 heartbeat 策略。
func (HeartbeatPolicy) FailureTimeout ¶
func (policy HeartbeatPolicy) FailureTimeout() time.Duration
FailureTimeout 返回从最近一次收到对端存活信号到判死的总超时窗口。
func (HeartbeatPolicy) NormalizeAndValidate ¶
func (policy HeartbeatPolicy) NormalizeAndValidate() (HeartbeatPolicy, error)
NormalizeAndValidate 归一化并校验 heartbeat 策略。
type HeartbeatStatus ¶
type HeartbeatStatus struct {
StartedAt time.Time
LastSentAt time.Time
LastReceivedAt time.Time
NextSendAt time.Time
FailureDeadline time.Time
ConsecutiveMisses int
Dead bool
}
HeartbeatStatus 描述某一时刻 heartbeat 的运行快照。
type IdleTunnelEviction ¶
type IdleTunnelEviction struct {
TunnelID string
Reason IdleTunnelEvictionReason
Cause error
}
IdleTunnelEviction 记录一次 idle tunnel 剔除结果。
type IdleTunnelEvictionReason ¶
type IdleTunnelEvictionReason string
IdleTunnelEvictionReason 描述 idle tunnel 被剔除的原因。
const ( // IdleTunnelEvictionReasonTTLExpired 表示 idle tunnel 超过 TTL。 IdleTunnelEvictionReasonTTLExpired IdleTunnelEvictionReason = "ttl_expired" // IdleTunnelEvictionReasonMissingInsertedAt 表示 idle tunnel 缺失入池时间。 IdleTunnelEvictionReasonMissingInsertedAt IdleTunnelEvictionReason = "missing_inserted_at" // IdleTunnelEvictionReasonBrokenState 表示 idle tunnel 状态已为 broken。 IdleTunnelEvictionReasonBrokenState IdleTunnelEvictionReason = "broken_state" // IdleTunnelEvictionReasonStateMismatch 表示 idle 池中的 tunnel 状态不再是 idle。 IdleTunnelEvictionReasonStateMismatch IdleTunnelEvictionReason = "state_mismatch" // IdleTunnelEvictionReasonDoneClosed 表示 tunnel 生命周期已结束。 IdleTunnelEvictionReasonDoneClosed IdleTunnelEvictionReason = "done_closed" // IdleTunnelEvictionReasonRuntimeError 表示 tunnel 已暴露运行时错误。 IdleTunnelEvictionReasonRuntimeError IdleTunnelEvictionReason = "runtime_error" // IdleTunnelEvictionReasonProbeFailed 表示 tunnel 探活失败。 IdleTunnelEvictionReasonProbeFailed IdleTunnelEvictionReason = "probe_failed" )
type InMemorySession ¶
type InMemorySession struct {
// contains filtered or unexported fields
}
InMemorySession 是线程安全的 Session 基础实现。
func NewInMemorySession ¶
func NewInMemorySession(meta SessionMeta, bindingInfo BindingInfo, capabilities SessionCapabilities) *InMemorySession
NewInMemorySession 创建内存版 Session 实例。
func (*InMemorySession) BindingInfo ¶
func (session *InMemorySession) BindingInfo() BindingInfo
BindingInfo 返回 binding 元信息。
func (*InMemorySession) Close ¶
func (session *InMemorySession) Close(ctx context.Context, reason error) error
Close 关闭 session 并触发 Done 信号。
func (*InMemorySession) Control ¶
func (session *InMemorySession) Control() (ControlChannel, error)
Control 返回 ControlChannel 能力。
func (*InMemorySession) Done ¶
func (session *InMemorySession) Done() <-chan struct{}
Done 返回会话结束通知。
func (*InMemorySession) MarkFailed ¶
func (session *InMemorySession) MarkFailed(cause error, retrying bool, giveUp bool) error
MarkFailed 将 session 标记为失败并记录映射上下文。
func (*InMemorySession) Meta ¶
func (session *InMemorySession) Meta() SessionMeta
Meta 返回 session 元信息快照。
func (*InMemorySession) Open ¶
func (session *InMemorySession) Open(ctx context.Context) error
Open 执行最小会话打开流程。
func (*InMemorySession) ProtocolState ¶
func (session *InMemorySession) ProtocolState() ProtocolState
ProtocolState 返回对外协议态视图。
func (*InMemorySession) State ¶
func (session *InMemorySession) State() SessionState
State 返回 transport 内部状态。
func (*InMemorySession) TunnelAcceptor ¶
func (session *InMemorySession) TunnelAcceptor() (TunnelAcceptor, error)
TunnelAcceptor 返回 TunnelAcceptor 能力。
func (*InMemorySession) TunnelPool ¶
func (session *InMemorySession) TunnelPool() (TunnelPool, error)
TunnelPool 返回 TunnelPool 能力。
func (*InMemorySession) TunnelProducer ¶
func (session *InMemorySession) TunnelProducer() (TunnelProducer, error)
TunnelProducer 返回 TunnelProducer 能力。
type InMemoryTunnelPool ¶
type InMemoryTunnelPool struct {
// contains filtered or unexported fields
}
InMemoryTunnelPool 是线程安全的内存实现。
func NewInMemoryTunnelPool ¶
func NewInMemoryTunnelPool() *InMemoryTunnelPool
NewInMemoryTunnelPool 创建默认的内存 tunnel 池。
func NewInMemoryTunnelPoolWithConfig ¶
func NewInMemoryTunnelPoolWithConfig(config TunnelPoolConfig) *InMemoryTunnelPool
NewInMemoryTunnelPoolWithConfig 使用指定配置创建内存 tunnel 池。
func (*InMemoryTunnelPool) Acquire ¶
func (pool *InMemoryTunnelPool) Acquire(ctx context.Context) (Tunnel, error)
Acquire 获取一条可用 idle tunnel,必要时阻塞等待。
func (*InMemoryTunnelPool) ClosedCount ¶
func (pool *InMemoryTunnelPool) ClosedCount() int
ClosedCount 返回累计降级关闭次数。
func (*InMemoryTunnelPool) Config ¶
func (pool *InMemoryTunnelPool) Config() TunnelPoolConfig
Config 返回当前池配置快照。
func (*InMemoryTunnelPool) EvictExpiredIdle ¶
func (pool *InMemoryTunnelPool) EvictExpiredIdle(now time.Time) []string
EvictExpiredIdle 按 idle_tunnel_ttl 清理过期空闲 tunnel。
func (*InMemoryTunnelPool) EvictZombieIdle ¶
func (pool *InMemoryTunnelPool) EvictZombieIdle(ctx context.Context, now time.Time) []IdleTunnelEviction
EvictZombieIdle 清理 idle 池中的僵尸 tunnel(状态异常、错误、探活失败、TTL 过期)。
func (*InMemoryTunnelPool) IdleCount ¶
func (pool *InMemoryTunnelPool) IdleCount() int
IdleCount 返回当前 idle tunnel 数量。
func (*InMemoryTunnelPool) InUseCount ¶
func (pool *InMemoryTunnelPool) InUseCount() int
InUseCount 返回当前 in-use tunnel 数量。
func (*InMemoryTunnelPool) PutIdle ¶
func (pool *InMemoryTunnelPool) PutIdle(tunnel Tunnel) error
PutIdle 将 tunnel 放入空闲池。
func (*InMemoryTunnelPool) Recycle ¶
func (pool *InMemoryTunnelPool) Recycle(ctx context.Context, tunnel Tunnel) (RecycleResult, error)
Recycle 把 in-use tunnel 回收入 idle 池;不满足条件时降级为关闭。
func (*InMemoryTunnelPool) RecycledCount ¶
func (pool *InMemoryTunnelPool) RecycledCount() int
RecycledCount 返回累计回收成功次数。
func (*InMemoryTunnelPool) Remove ¶
func (pool *InMemoryTunnelPool) Remove(tunnelID string) error
Remove 从池中移除指定 tunnel。
type KeepalivePolicy ¶
type KeepalivePolicy struct {
// IdleTTL 表示 idle tunnel 在池中的最大存活时间;0 表示不启用 TTL 轮换。
IdleTTL time.Duration
// ProbeInterval 表示应用层探活间隔;0 表示依赖传输层 keepalive。
ProbeInterval time.Duration
// ProbeTimeout 表示单次探活超时。
ProbeTimeout time.Duration
// ProbeMaxFailures 表示连续探活失败阈值。
ProbeMaxFailures int
}
KeepalivePolicy 描述 binding 对 idle tunnel 的保活治理参数。
func DefaultKeepalivePolicyForBinding ¶
func DefaultKeepalivePolicyForBinding(bindingType BindingType) KeepalivePolicy
DefaultKeepalivePolicyForBinding 返回 binding 的默认保活治理参数。
type PoolDiagnosticFields ¶
type PoolDiagnosticFields struct {
SessionID string
SessionEpoch uint64
Binding string
IdleCount int
InUseCount int
TargetIdleCount int
LastError string
}
PoolDiagnosticFields 描述 tunnel pool 维度的结构化日志字段。
func BuildPoolDiagnosticFields ¶
func BuildPoolDiagnosticFields(report TunnelPoolReport, bindingInfo BindingInfo, lastError error) PoolDiagnosticFields
BuildPoolDiagnosticFields 构造 pool report 维度统一诊断字段。
func (PoolDiagnosticFields) Fields ¶
func (fields PoolDiagnosticFields) Fields() map[string]any
Fields 返回适合结构化日志输出的字段映射。
type PoolRefillDiagnosticFields ¶
type PoolRefillDiagnosticFields struct {
SessionID string
SessionEpoch uint64
Binding string
RequestID string
Reason TunnelRefillReason
RequestedIdleDelta int
EffectiveTarget int
OpenedCount int
FailedCount int
AfterIdleCount int
LastError string
}
PoolRefillDiagnosticFields 描述补池动作维度的结构化日志字段。
func BuildPoolRefillDiagnosticFields ¶
func BuildPoolRefillDiagnosticFields( request TunnelRefillRequest, result RefillResult, bindingInfo BindingInfo, lastError error, ) PoolRefillDiagnosticFields
BuildPoolRefillDiagnosticFields 构造补池动作统一诊断字段。
func (PoolRefillDiagnosticFields) Fields ¶
func (fields PoolRefillDiagnosticFields) Fields() map[string]any
Fields 返回适合结构化日志输出的字段映射。
type PrioritizedControlChannel ¶
type PrioritizedControlChannel interface {
ControlChannel
WritePrioritizedControlFrame(ctx context.Context, frame PrioritizedControlFrame) error
}
PrioritizedControlChannel 为支持优先级发送调度的可选控制面能力。
type PrioritizedControlFrame ¶
type PrioritizedControlFrame struct {
Priority ControlMessagePriority
Frame ControlFrame
}
PrioritizedControlFrame 描述带优先级的控制帧。
type PriorityControlQueue ¶
type PriorityControlQueue struct {
// contains filtered or unexported fields
}
PriorityControlQueue 提供控制面发送队列的最小优先级调度实现。
func NewPriorityControlQueue ¶
func NewPriorityControlQueue() *PriorityControlQueue
NewPriorityControlQueue 创建优先级控制队列。
func (*PriorityControlQueue) Dequeue ¶
func (queue *PriorityControlQueue) Dequeue() (PrioritizedControlFrame, bool)
Dequeue 按 high -> normal -> low 的顺序出队。
func (*PriorityControlQueue) Enqueue ¶
func (queue *PriorityControlQueue) Enqueue(frame PrioritizedControlFrame)
Enqueue 按优先级将控制帧入队。
type ProtocolState ¶
type ProtocolState string
ProtocolState 描述对外暴露的 LTFP 协议态。
const ( // ProtocolStateConnecting 表示正在连接或重连中。 ProtocolStateConnecting ProtocolState = "CONNECTING" // ProtocolStateAuthenticating 表示控制面初始化/认证阶段。 ProtocolStateAuthenticating ProtocolState = "AUTHENTICATING" // ProtocolStateActive 表示会话处于可服务状态。 ProtocolStateActive ProtocolState = "ACTIVE" // ProtocolStateDraining 表示会话排空中。 ProtocolStateDraining ProtocolState = "DRAINING" // ProtocolStateStale 表示已激活会话变为失效状态。 ProtocolStateStale ProtocolState = "STALE" // ProtocolStateClosed 表示会话已关闭。 ProtocolStateClosed ProtocolState = "CLOSED" )
func MapSessionStateToProtocolState ¶
func MapSessionStateToProtocolState(sessionState SessionState, mapping FailedMappingContext) ProtocolState
MapSessionStateToProtocolState 将 transport 内部态映射为对外协议态。
type RecycleResult ¶
type RecycleResult string
RecycleResult 描述一次回收操作的最终结果。
const ( // RecycleResultRecycled 表示回收成功并重新回到 idle 池。 RecycleResultRecycled RecycleResult = "recycled" // RecycleResultClosed 表示回收失败并降级关闭。 RecycleResultClosed RecycleResult = "closed" )
type RefillController ¶
type RefillController struct {
// contains filtered or unexported fields
}
RefillController 管理补池、限流和去重逻辑。
func NewRefillController ¶
func NewRefillController(pool TunnelPool, producer TunnelProducer, config RefillControllerConfig) (*RefillController, error)
NewRefillController 创建补池控制器。
func (*RefillController) HandleRefillRequest ¶
func (controller *RefillController) HandleRefillRequest(ctx context.Context, request TunnelRefillRequest) (RefillResult, error)
HandleRefillRequest 处理控制面补池请求(包含幂等去重)。
func (*RefillController) RefillToTarget ¶
func (controller *RefillController) RefillToTarget(ctx context.Context, targetIdle int) (RefillResult, error)
RefillToTarget 将空闲池平滑补充到目标容量。
type RefillControllerConfig ¶
type RefillControllerConfig struct {
MinIdleTunnels int
MaxIdleTunnels int
MaxInFlightTunnelOpens int
TunnelOpenRateLimit rate.Limit
TunnelOpenBurst int
RequestDeduplicateTTL time.Duration
}
RefillControllerConfig 描述补池控制循环配置。
func DefaultRefillControllerConfig ¶
func DefaultRefillControllerConfig() RefillControllerConfig
DefaultRefillControllerConfig 返回默认补池配置。
func (RefillControllerConfig) NormalizeAndValidate ¶
func (config RefillControllerConfig) NormalizeAndValidate() (RefillControllerConfig, error)
NormalizeAndValidate 归一化并校验补池配置。
type RefillRequestDeduplicator ¶
type RefillRequestDeduplicator struct {
// contains filtered or unexported fields
}
RefillRequestDeduplicator 维护补池请求去重窗口。
func NewRefillRequestDeduplicator ¶
func NewRefillRequestDeduplicator(ttl time.Duration) *RefillRequestDeduplicator
NewRefillRequestDeduplicator 创建请求去重器。
type RefillResult ¶
type RefillResult struct {
RequestID string
Deduplicated bool
BeforeIdleCount int
RequestedTargetIdle int
EffectiveTargetIdle int
OpenedCount int
FailedCount int
AfterIdleCount int
FirstFailure error
}
RefillResult 描述一次补池动作结果。
type Session ¶
type Session interface {
ID() string
Meta() SessionMeta
State() SessionState
BindingInfo() BindingInfo
ProtocolState() ProtocolState
Open(ctx context.Context) error
Close(ctx context.Context, reason error) error
MarkFailed(cause error, retrying bool, giveUp bool) error
Control() (ControlChannel, error)
TunnelProducer() (TunnelProducer, error)
TunnelAcceptor() (TunnelAcceptor, error)
TunnelPool() (TunnelPool, error)
Done() <-chan struct{}
Err() error
}
Session 定义 transport 聚合根接口。
type SessionCapabilities ¶
type SessionCapabilities struct {
ControlChannel ControlChannel
Producer TunnelProducer
Acceptor TunnelAcceptor
Pool TunnelPool
}
SessionCapabilities 统一封装 session 可用能力。
type SessionDiagnosticFields ¶
type SessionDiagnosticFields struct {
SessionID string
SessionEpoch uint64
Binding string
State SessionState
ProtocolState ProtocolState
LastError string
}
SessionDiagnosticFields 描述 session 维度的结构化日志字段。
func BuildSessionDiagnosticFields ¶
func BuildSessionDiagnosticFields( meta SessionMeta, bindingInfo BindingInfo, state SessionState, protocolState ProtocolState, lastError error, ) SessionDiagnosticFields
BuildSessionDiagnosticFields 构造 session 维度统一诊断字段。
func BuildSessionDiagnosticFieldsForSession ¶
func BuildSessionDiagnosticFieldsForSession(session Session) SessionDiagnosticFields
BuildSessionDiagnosticFieldsForSession 从 Session 直接提取诊断字段。
func (SessionDiagnosticFields) Fields ¶
func (fields SessionDiagnosticFields) Fields() map[string]any
Fields 返回适合结构化日志输出的字段映射。
type SessionMeta ¶
type SessionMeta struct {
SessionID string
SessionEpoch uint64
NodeID string
Labels map[string]string
LastError string
}
SessionMeta 描述 session 级元数据。
type SessionState ¶
type SessionState string
SessionState 描述 transport 层内部会话状态。
const ( // SessionStateIdle 表示会话尚未打开。 SessionStateIdle SessionState = "idle" // SessionStateConnecting 表示正在建立底层连接。 SessionStateConnecting SessionState = "connecting" // SessionStateConnected 表示底层连接已建立。 SessionStateConnected SessionState = "connected" // SessionStateControlReady 表示控制面通道已可用。 SessionStateControlReady SessionState = "control_ready" // SessionStateAuthenticated 表示认证完成且可处理业务控制面消息。 SessionStateAuthenticated SessionState = "authenticated" // SessionStateDraining 表示会话进入排空阶段。 SessionStateDraining SessionState = "draining" // SessionStateFailed 表示会话发生内部失败。 SessionStateFailed SessionState = "failed" // SessionStateClosed 表示会话已关闭。 SessionStateClosed SessionState = "closed" )
func (SessionState) IsTerminal ¶
func (sessionState SessionState) IsTerminal() bool
IsTerminal 判断会话状态是否为终止态。
type TransportMetricsRecorder ¶
type TransportMetricsRecorder struct {
// contains filtered or unexported fields
}
TransportMetricsRecorder 维护 transport 关键指标的并发安全聚合。
func NewTransportMetricsRecorder ¶
func NewTransportMetricsRecorder() *TransportMetricsRecorder
NewTransportMetricsRecorder 创建默认指标记录器。
func (*TransportMetricsRecorder) IncBrokenTunnel ¶
func (recorder *TransportMetricsRecorder) IncBrokenTunnel()
IncBrokenTunnel 记录一次 broken tunnel 事件。
func (*TransportMetricsRecorder) IncOpenTimeout ¶
func (recorder *TransportMetricsRecorder) IncOpenTimeout()
IncOpenTimeout 记录一次 open timeout。
func (*TransportMetricsRecorder) IncReset ¶
func (recorder *TransportMetricsRecorder) IncReset()
IncReset 记录一次 reset 事件。
func (*TransportMetricsRecorder) ObserveHeartbeatRTT ¶
func (recorder *TransportMetricsRecorder) ObserveHeartbeatRTT(roundTripTime time.Duration)
ObserveHeartbeatRTT 记录最近一次 heartbeat RTT。
func (*TransportMetricsRecorder) ObservePoolCounts ¶
func (recorder *TransportMetricsRecorder) ObservePoolCounts(idleCount int, inUseCount int)
ObservePoolCounts 记录当前 pool 的 idle/in-use 数量。
func (*TransportMetricsRecorder) ObserveRefill ¶
func (recorder *TransportMetricsRecorder) ObserveRefill(openedCount int, observedAt time.Time)
ObserveRefill 记录一次补池结果,用于计算补池速率。
func (*TransportMetricsRecorder) Snapshot ¶
func (recorder *TransportMetricsRecorder) Snapshot(now time.Time) TransportMetricsSnapshot
Snapshot 返回当前指标快照。
type TransportMetricsSnapshot ¶
type TransportMetricsSnapshot struct {
HeartbeatRTT time.Duration
IdleCount int
InUseCount int
RefillRatePerSec float64
RefillOpenedTotal uint64
OpenTimeoutCount uint64
ResetCount uint64
BrokenCount uint64
}
TransportMetricsSnapshot 描述 transport 关键指标快照。
func (TransportMetricsSnapshot) Fields ¶
func (snapshot TransportMetricsSnapshot) Fields() map[string]any
Fields 返回可直接用于结构化日志/指标导出的字段映射。
type Tunnel ¶
type Tunnel interface {
io.Reader
io.Writer
io.Closer
ID() string
Meta() TunnelMeta
State() TunnelState
BindingInfo() BindingInfo
CloseWrite() error
Reset(cause error) error
SetDeadline(deadline time.Time) error
SetReadDeadline(deadline time.Time) error
SetWriteDeadline(deadline time.Time) error
// Flush 在回收前尝试清理本地读缓存;若存在脏数据应返回错误。
Flush() error
// ReuseCount 返回当前 tunnel 已完成的回收轮次(0-based)。
ReuseCount() int
// Recyclable 报告当前 tunnel 是否满足回收前提。
Recyclable() bool
Done() <-chan struct{}
Err() error
}
Tunnel 定义数据面最小字节流抽象。
type TunnelAcceptor ¶
TunnelAcceptor 定义 Server 侧接收 tunnel 能力。
type TunnelDiagnosticFields ¶
type TunnelDiagnosticFields struct {
SessionID string
SessionEpoch uint64
Binding string
TunnelID string
TunnelState TunnelState
LastError string
}
TunnelDiagnosticFields 描述 tunnel 维度的结构化日志字段。
func BuildTunnelDiagnosticFields ¶
func BuildTunnelDiagnosticFields( meta TunnelMeta, bindingInfo BindingInfo, state TunnelState, lastError error, ) TunnelDiagnosticFields
BuildTunnelDiagnosticFields 构造 tunnel 维度统一诊断字段。
func BuildTunnelDiagnosticFieldsForTunnel ¶
func BuildTunnelDiagnosticFieldsForTunnel(tunnel Tunnel) TunnelDiagnosticFields
BuildTunnelDiagnosticFieldsForTunnel 从 Tunnel 直接提取诊断字段。
func (TunnelDiagnosticFields) Fields ¶
func (fields TunnelDiagnosticFields) Fields() map[string]any
Fields 返回适合结构化日志输出的字段映射。
type TunnelHealthProber ¶
TunnelHealthProber 定义 tunnel 可选探活能力。
type TunnelMeta ¶
type TunnelMeta struct {
TunnelID string
SessionID string
SessionEpoch uint64
CreatedAt time.Time
Labels map[string]string
}
TunnelMeta 描述 tunnel 级元数据。
type TunnelPool ¶
type TunnelPool interface {
PutIdle(tunnel Tunnel) error
Acquire(ctx context.Context) (Tunnel, error)
Remove(tunnelID string) error
Recycle(ctx context.Context, tunnel Tunnel) (RecycleResult, error)
IdleCount() int
InUseCount() int
RecycledCount() int
ClosedCount() int
}
TunnelPool 定义 tunnel 池最小接口。
type TunnelPoolConfig ¶
type TunnelPoolConfig struct {
MinIdleTunnels int
MaxIdleTunnels int
IdleTunnelTTL time.Duration
AcquireTimeout time.Duration
}
TunnelPoolConfig 描述 tunnel 池容量和超时参数。
func DefaultTunnelPoolConfig ¶
func DefaultTunnelPoolConfig() TunnelPoolConfig
DefaultTunnelPoolConfig 返回默认池配置。
func (TunnelPoolConfig) NormalizeAndValidate ¶
func (config TunnelPoolConfig) NormalizeAndValidate() (TunnelPoolConfig, error)
NormalizeAndValidate 归一化并校验配置。
type TunnelPoolReport ¶
type TunnelPoolReport struct {
SessionID string
SessionEpoch uint64
IdleCount int
InUseCount int
TargetIdleCount int
Timestamp time.Time
}
TunnelPoolReport 描述当前会话的池状态上报。
func (TunnelPoolReport) Validate ¶
func (report TunnelPoolReport) Validate() error
Validate 校验 TunnelPoolReport 字段合法性。
type TunnelProducer ¶
TunnelProducer 定义 Agent 侧主动建 tunnel 能力。
type TunnelRefillReason ¶
type TunnelRefillReason string
TunnelRefillReason 表示补池触发原因。
const ( // TunnelRefillReasonLowWatermark 表示 idle 低于水位线。 TunnelRefillReasonLowWatermark TunnelRefillReason = "low_watermark" // TunnelRefillReasonAcquireTimeout 表示 acquire 超时触发补池。 TunnelRefillReasonAcquireTimeout TunnelRefillReason = "acquire_timeout" // TunnelRefillReasonStartup 表示启动阶段预热补池。 TunnelRefillReasonStartup TunnelRefillReason = "startup" // TunnelRefillReasonManual 表示人工触发补池。 TunnelRefillReasonManual TunnelRefillReason = "manual" )
type TunnelRefillRequest ¶
type TunnelRefillRequest struct {
SessionID string
SessionEpoch uint64
RequestID string
RequestedIdleDelta int
Reason TunnelRefillReason
Timestamp time.Time
}
TunnelRefillRequest 描述控制面补池请求。
func (TunnelRefillRequest) Normalize ¶
func (request TunnelRefillRequest) Normalize() TunnelRefillRequest
Normalize 归一化补池请求字段。
func (TunnelRefillRequest) Validate ¶
func (request TunnelRefillRequest) Validate() error
Validate 校验补池请求是否合法。
type TunnelState ¶
type TunnelState string
TunnelState 描述 transport 层内部 tunnel 状态。
const ( // TunnelStateOpening 表示 tunnel 正在建立。 TunnelStateOpening TunnelState = "opening" // TunnelStateIdle 表示 tunnel 空闲可分配。 TunnelStateIdle TunnelState = "idle" // TunnelStateReserved 表示 tunnel 已预留但尚未进入 active。 TunnelStateReserved TunnelState = "reserved" // TunnelStateActive 表示 tunnel 正在承载流量。 TunnelStateActive TunnelState = "active" // TunnelStateClosing 表示 tunnel 正在关闭。 TunnelStateClosing TunnelState = "closing" // TunnelStateClosed 表示 tunnel 已正常关闭。 TunnelStateClosed TunnelState = "closed" // TunnelStateBroken 表示 tunnel 异常损坏不可复用。 TunnelStateBroken TunnelState = "broken" )
func (TunnelState) IsTerminal ¶
func (tunnelState TunnelState) IsTerminal() bool
IsTerminal 判断 tunnel 状态是否为终止态。
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package grpcbinding 提供 grpc_h2 传输绑定实现。
|
Package grpcbinding 提供 grpc_h2 传输绑定实现。 |
|
Package quicbinding 提供 quic_native 传输绑定实现。
|
Package quicbinding 提供 quic_native 传输绑定实现。 |
|
Package tcpbinding 提供 tcp_framed 传输绑定实现。
|
Package tcpbinding 提供 tcp_framed 传输绑定实现。 |