p2p

package
v2.2.0 Latest Latest
Warning

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

Go to latest
Published: Jul 7, 2026 License: Apache-2.0 Imports: 17 Imported by: 0

Documentation

Overview

Package p2p 提供基于 ICE/STUN/TURN 的点对点连接能力。

ICE 打通 UDP 路径后在其上承载 QUIC(可靠 + 多路 + TLS 1.3), 信令复用 tunnel gateway 控制流;TURN 由外部 coturn 提供。

详见 docs/design-p2p.md。

Index

Constants

View Source
const (
	DefaultSTUNURL      = "stun:118.178.168.253:3478"
	DefaultTURNURL      = "turn:118.178.168.253:3478"
	DefaultTURNUsername = "lava"
	DefaultTURNPassword = "lava-p2p-dev"
	DefaultTURNRealm    = "lava.dev"
)

默认 dev 服务器 coturn(见 deploy/coturn/turnserver.conf)。

Variables

View Source
var (
	ErrSignalingClosed = errors.New("p2p: signaling closed")
	ErrPeerNotFound    = errors.New("p2p: peer not found")
	ErrICEFailed       = errors.New("p2p: ice connection failed")
	ErrTimeout         = errors.New("p2p: timeout")
)

Functions

func AuthTokenFromEnv

func AuthTokenFromEnv() string

AuthTokenFromEnv 读取 P2P 信令鉴权 token。

func PeerIDFromEnv

func PeerIDFromEnv() string

PeerIDFromEnv 读取 P2P 节点 ID;未配置时返回空字符串(表示不启用 P2P)。

Types

type CandidatePairInfo

type CandidatePairInfo struct {
	LocalType  string `json:"local_type"`
	RemoteType string `json:"remote_type"`
	LocalAddr  string `json:"local_addr"`
	RemoteAddr string `json:"remote_addr"`
}

CandidatePairInfo 是 ICE 选路结果摘要。

type Client

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

Client 是 Coordinator 之上的高层封装:按 peerID 复用 P2P 连接、 在连接失效时自动重连,并提供 net.Conn / gRPC / HTTP 适配,便于集成。

典型用法:

client := p2p.NewClient(coord)
conn, err := client.OpenStream(ctx, "node-b") // 返回 net.Conn

// gRPC over P2P:
cc, _ := grpc.NewClient("passthrough:///node-b",
    grpc.WithContextDialer(client.NetDialer()),
    grpc.WithTransportCredentials(insecure.NewCredentials()))

// HTTP over P2P(host 即 peerID):
resp, _ := client.HTTPClient().Get("http://node-b/healthz")

func NewClient

func NewClient(coord Coordinator) *Client

NewClient 基于已有 Coordinator 创建池化客户端。

func (*Client) Close

func (c *Client) Close() error

Close 关闭所有池化连接(不关闭底层 Coordinator)。

func (*Client) Conn

func (c *Client) Conn(ctx context.Context, peerID string) (PeerConn, error)

Conn 返回到 peerID 的一条健康连接:命中缓存则复用, 连接已关闭则自动重连,无连接则首次拨号。并发调用同一 peer 会串行化。

func (*Client) HTTPClient

func (c *Client) HTTPClient() *http.Client

HTTPClient 返回经 P2P 拨号的 http.Client(URL host 即 peerID)。

func (*Client) HTTPTransport

func (c *Client) HTTPTransport() *http.Transport

HTTPTransport 返回经 P2P 拨号的 http.Transport(URL host 即 peerID)。

func (*Client) NetDialer

func (c *Client) NetDialer() func(context.Context, string) (net.Conn, error)

NetDialer 返回可用于 grpc.WithContextDialer 的拨号函数(addr 即 peerID)。

func (*Client) OpenStream

func (c *Client) OpenStream(ctx context.Context, peerID string) (net.Conn, error)

OpenStream 在到 peerID 的复用连接上打开一条新流(即 net.Conn)。 若底层连接在打开时已失效,会重连一次后重试。

type Config

type Config struct {
	// STUNURLs STUN 服务器列表。
	STUNURLs []string `yaml:"stun_urls"`
	// TURN 配置(外部 coturn)。
	TURN TURNConfig `yaml:"turn"`
	// ICE 超时。
	ICETimeout time.Duration `yaml:"ice_timeout"`
	// SignalingAddr tunnel gateway 信令地址(P1 后接入 tunnel 控制流)。
	SignalingAddr string `yaml:"signaling_addr"`
	// AuthToken 注册/建连鉴权 token,对应 tunnel TokenAuthProvider。
	AuthToken string `yaml:"auth_token"`
	// Insecure 跳过 QUIC TLS 校验(仅开发)。
	Insecure bool `yaml:"insecure"`
	// TLS QUIC 证书(生产环境);设置 cert/key 后自动启用 TLS。
	TLS TLSFileConfig `yaml:"tls"`
	// Reconnect 断线重连策略(Reconnect 使用;Dial 仍为单次尝试)。
	Reconnect ReconnectConfig `yaml:"reconnect"`
}

Config P2P 模块配置。

func ConfigFromEnv

func ConfigFromEnv() Config

ConfigFromEnv 从环境变量加载 P2P 配置(在 DefaultConfig 基础上覆盖)。

环境变量:

  • P2P_STUN_URLS 逗号分隔 STUN URL
  • P2P_TURN_URL TURN URL(设为空或 P2P_TURN_DISABLED=1 禁用)
  • P2P_TURN_DISABLED true/1 禁用 TURN
  • P2P_TURN_USER TURN 用户名(HMAC 模式下作为 user_id 后缀,缺省用 peerID)
  • P2P_TURN_PASS TURN 静态密码(未设置 P2P_TURN_SECRET 时使用)
  • P2P_TURN_SECRET coturn static-auth-secret,启用 HMAC 临时凭证
  • P2P_TURN_CRED_TTL 临时凭证 TTL(如 24h,默认 24h)
  • P2P_RECONNECT_MAX_ATTEMPTS Reconnect 最大尝试次数(默认 3)
  • P2P_RECONNECT_BACKOFF Reconnect 重试间隔(如 1s)
  • P2P_INSECURE true/1 跳过 QUIC TLS 校验(仅开发)
  • P2P_CERT_FILE QUIC TLS 证书
  • P2P_KEY_FILE QUIC TLS 私钥
  • P2P_CA_FILE QUIC TLS CA(可选)
  • TUNNEL_AUTH_TOKEN 信令鉴权 token(若 P2P_AUTH_TOKEN 未设置)
  • P2P_AUTH_TOKEN 信令鉴权 token

func DefaultConfig

func DefaultConfig() Config

DefaultConfig 返回面向 dev 环境的默认配置。

func (Config) ICEURLs

func (c Config) ICEURLs() []string

ICEURLs 返回传给 pion/ice 的 STUN/TURN URL 列表(含 TURN 凭证)。

func (Config) TransportOptions

func (c Config) TransportOptions() *tunnel.TransportOptions

TransportOptions 转为 tunnel QUIC 传输选项。

type ConnectionStats

type ConnectionStats struct {
	RemotePeerID      string            `json:"remote_peer_id"`
	Pair              CandidatePairInfo `json:"pair"`
	ConnectedAt       time.Time         `json:"connected_at"`
	ConnectDurationMs int64             `json:"connect_duration_ms"`
	RTTMs             float64           `json:"rtt_ms,omitempty"`
	PairState         string            `json:"pair_state,omitempty"`
	Nominated         bool              `json:"nominated,omitempty"`
}

ConnectionStats 单条 P2P 连接的 ICE 选路摘要。

type Coordinator

type Coordinator interface {
	Listen(ctx context.Context, selfID string) (Listener, error)
	Dial(ctx context.Context, peerID string) (PeerConn, error)
	// Reconnect 关闭与 peer 的现有连接并重新 ICE+QUIC 协商(带退避重试)。
	Reconnect(ctx context.Context, peerID string) (PeerConn, error)
	Close() error
	Stats() Stats
}

Coordinator 管理本节点的 P2P 连接(Listen / Dial / Reconnect)。

func NewCoordinator

func NewCoordinator(cfg Config, broker signaling.Broker, selfID string) Coordinator

NewCoordinator 创建 P2P 协调器。

func NewCoordinatorWithMetrics

func NewCoordinatorWithMetrics(cfg Config, broker signaling.Broker, selfID string, met *MetricsRecorder) Coordinator

NewCoordinatorWithMetrics 创建 P2P 协调器并绑定可选 metrics。

type Listener

type Listener interface {
	Accept(ctx context.Context) (PeerConn, error)
	Close() error
}

Listener 接受来自对端的 P2P 入站连接。

type MetricsRecorder

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

MetricsRecorder 上报 P2P 建连指标(scope 为 nil 时为 noop)。

func NewMetricsRecorder

func NewMetricsRecorder(scope metrics.Metric) *MetricsRecorder

NewMetricsRecorder 创建指标记录器;scope 为 nil 时返回 nil(noop)。

func (*MetricsRecorder) ObserveConnect

func (m *MetricsRecorder) ObserveConnect(pair CandidatePairInfo, d time.Duration)

ObserveConnect 记录一次成功建连。

func (*MetricsRecorder) ObserveDialFailure

func (m *MetricsRecorder) ObserveDialFailure()

ObserveDialFailure 记录一次拨号失败。

func (*MetricsRecorder) SetActiveConnections

func (m *MetricsRecorder) SetActiveConnections(active, relay int)

SetActiveConnections 更新当前活跃连接数与中继连接数。

type PacketConn

type PacketConn interface {
	net.PacketConn
	io.Closer
}

PacketConn 是 ICE 打通后可用于 QUIC 的 UDP 数据报连接。

type PeerConn

type PeerConn interface {
	tunnel.Session
	RemotePeerID() string
	LocalPeerID() string
	// SelectedPair 返回 ICE 选路摘要(host/srflx/relay + 地址),用于可观测。
	SelectedPair() CandidatePairInfo
}

PeerConn 是一条已建立的 P2P 连接,实现 tunnel.Session 语义。

type ReconnectConfig

type ReconnectConfig struct {
	// MaxAttempts 最大尝试次数(含首次),默认 3。
	MaxAttempts int `yaml:"max_attempts"`
	// Backoff 重试间隔,默认 1s。
	Backoff time.Duration `yaml:"backoff"`
}

ReconnectConfig 断线重连退避策略。

type Role

type Role int

Role 表示 ICE 协商中的角色。

const (
	RoleDialer Role = iota
	RoleListener
)

type Stats

type Stats struct {
	SelfID            string            `json:"self_id"`
	Listening         bool              `json:"listening"`
	ActiveConnections int               `json:"active_connections"`
	RelayConnections  int               `json:"relay_connections"`
	Connections       []ConnectionStats `json:"connections"`
}

Stats P2P 节点摘要。

type TLSFileConfig

type TLSFileConfig struct {
	CertFile string `yaml:"cert_file"`
	KeyFile  string `yaml:"key_file"`
	CAFile   string `yaml:"ca_file"`
}

TLSFileConfig QUIC TLS 文件路径。

type TURNConfig

type TURNConfig struct {
	URL      string `yaml:"url"`
	Username string `yaml:"username"`
	Password string `yaml:"password"`
	Realm    string `yaml:"realm"`
	// AuthSecret 与 coturn static-auth-secret 一致时,按 TURN REST API 生成临时凭证。
	AuthSecret string        `yaml:"auth_secret"`
	CredTTL    time.Duration `yaml:"cred_ttl"`
	// Disabled 为 true 时不使用 TURN(即使 URL 有默认值)。
	Disabled bool `yaml:"disabled"`
}

TURNConfig coturn 客户端配置。

func (TURNConfig) Resolved

func (t TURNConfig) Resolved(peerID string) (username, password string, err error)

Resolved 返回实际用于 ICE 的 TURN 凭证(含 HMAC 临时密码)。

Directories

Path Synopsis
P2P 示例:本地内存信令 / tunnel gateway 全链路 / dev coturn 验证。
P2P 示例:本地内存信令 / tunnel gateway 全链路 / dev coturn 验证。
Package p2pbuilder 装配 P2P Coordinator:tunnel agent 信令、lifecycle 与 debug 端点。
Package p2pbuilder 装配 P2P Coordinator:tunnel agent 信令、lifecycle 与 debug 端点。
tunnelsig
Package tunnelsig 通过 tunnel gateway 控制流交换 ICE 信令。
Package tunnelsig 通过 tunnel gateway 控制流交换 ICE 信令。
Package turncred 生成 coturn TURN REST API 临时凭证(HMAC-SHA1)。
Package turncred 生成 coturn TURN REST API 临时凭证(HMAC-SHA1)。

Jump to

Keyboard shortcuts

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