transport

package
v0.0.0-...-c041292 Latest Latest
Warning

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

Go to latest
Published: Mar 30, 2026 License: MIT Imports: 14 Imported by: 0

Documentation

Index

Constants

View Source
const (
	// ControlFrameTypeFragment 是控制面分块帧的保留类型。
	ControlFrameTypeFragment uint16 = 0xFFFF

	// DefaultControlFragmentMaxPayloadSize 是控制面分块的默认单帧上限。
	DefaultControlFragmentMaxPayloadSize = 64 * 1024
)
View Source
const (
	// ControlFrameTypeHeartbeatPing 表示控制面 heartbeat ping 帧。
	ControlFrameTypeHeartbeatPing uint16 = 0xFFF0
	// ControlFrameTypeHeartbeatPong 表示控制面 heartbeat pong 帧。
	ControlFrameTypeHeartbeatPong uint16 = 0xFFF1
)
View Source
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

View Source
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

func IsErrorKind(err error, kind ErrorKind) bool

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

type ControlFrame struct {
	Type    uint16
	Payload []byte
}

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

type Error struct {
	Kind      ErrorKind
	Op        string
	Message   string
	Temporary bool
	Cause     error
}

Error 定义 transport 统一错误结构。

func WrapError

func WrapError(kind ErrorKind, operation string, cause error, message string) *Error

WrapError 将底层错误包装为 transport Error。

func Wrapf

func Wrapf(kind ErrorKind, operation string, cause error, format string, arguments ...any) *Error

Wrapf 按格式化方式包装错误并保留底层 cause。

func (*Error) Error

func (err *Error) Error() string

Error 返回错误字符串。

func (*Error) Unwrap

func (err *Error) Unwrap() error

Unwrap 返回底层错误,支持 errors.Is / errors.As。

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) Err

func (session *InMemorySession) Err() error

Err 返回会话最近错误。

func (*InMemorySession) ID

func (session *InMemorySession) ID() string

ID 返回 sessionId。

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

Dequeue 按 high -> normal -> low 的顺序出队。

func (*PriorityControlQueue) Enqueue

func (queue *PriorityControlQueue) Enqueue(frame PrioritizedControlFrame)

Enqueue 按优先级将控制帧入队。

func (*PriorityControlQueue) Len

func (queue *PriorityControlQueue) Len() int

Len 返回当前队列总长度。

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 创建请求去重器。

func (*RefillRequestDeduplicator) MarkSeen

func (deduplicator *RefillRequestDeduplicator) MarkSeen(requestID string, now time.Time) bool

MarkSeen 标记 requestId 已处理,返回是否首次处理。

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

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

type TunnelAcceptor interface {
	AcceptTunnel(ctx context.Context) (Tunnel, error)
}

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

type TunnelHealthProber interface {
	Probe(ctx context.Context) error
}

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

type TunnelProducer interface {
	OpenTunnel(ctx context.Context) (Tunnel, error)
}

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 状态是否为终止态。

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 传输绑定实现。

Jump to

Keyboard shortcuts

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