SFU 服务设计与实现:基于 webrtc-rs 的 Rust 多方音视频转发
SFU 服务设计与实现:基于 webrtc-rs 的 Rust 多方音视频转发
在
mindit去中心化系统中,SFU 仅作为 RTP 数据包的“无情转发机器”。由于我们采用了 OpenMLS 实现了端到端加密 (E2EE),SFU 服务器即使被抓包或攻破,也完全无法解密用户的音视频流 Payload,从而在架构层实现了真正的零信任 (Zero-Trust)安全。本文是《Rust + OpenMLS 端到端加密音视频系统》系列的 SFU 篇。
系列概览:[架构总览] → [MLS 加密] → [SFU 实现] → [信令设计] → [本地加密存储]
目录
- 什么是 SFU,为什么选它
- 整体架构
- 核心数据结构设计
- PeerConnection 构建:网络配置详解
- Publish 流程
- Subscribe 流程
- 动态新成员加入:重协商机制
- RTP 转发与 PLI 自愈
- 带宽监控:无锁原子计数器
- 鉴权设计:防越权攻击
- 房间生命周期管理
- 管控能力:强踢与会话终止
- 生产环境部署要点
- 已知局限与后续规划
1. 什么是 SFU,为什么选它
多方音视频有三种经典拓扑:
| 拓扑 | 原理 | 优点 | 缺点 |
|---|---|---|---|
| Mesh | 每两个人直连 | 无服务器压力 | N 路上行,N=4 即崩溃 |
| MCU | 服务器混流后转发 | 客户端带宽最低 | 服务器 CPU 极高,破坏 E2EE |
| SFU | 服务器只转发 RTP,不解码 | 横向扩展好,保留 E2EE | 客户端多路下行 |
本项目选择 SFU,原因有两点:
- E2EE 兼容:SFU 只看 RTP Header(路由用),不碰 Payload(内容),不破坏 MLS 加密的语义安全性。
- Rust 原生:webrtc-rs 提供了完整的 WebRTC 协议栈,与 Axum + Tokio 异步生态无缝集成。
2. 整体架构
┌─────────────────────────────────────────────────────────────────┐
│ Client (Browser/Tauri) │
│ │
│ ┌──────────┐ HTTP POST /publish ┌──────────────────────┐ │
│ │Publisher │ ──────────────────────▶│ │ │
│ │(Offer) │ ◀────────────────────── │ SFU Server │ │
│ └──────────┘ Answer (base64 SDP) │ │ │
│ │ ┌────────────────┐ │ │
│ ┌──────────┐ HTTP POST /subscribe │ │ RoomRegistry │ │ │
│ │Subscriber│ ──────────────────────▶│ │ │ │ │
│ │(Offer) │ ◀────────────────────── │ │ publishers: {} │ │ │
│ └──────────┘ Answer (base64 SDP) │ │ subscribers:{} │ │ │
│ │ │ signal_members │ │ │
│ WebSocket 信令(Offer 重协商) │ └────────────────┘ │ │
│ ┌──────────┐ ◀──────────────────────│ │ │
│ │Old User │ ──────────────────────▶│ broadcast::channel │ │
│ └──────────┘ POST /subscribe/answer│ video_tx / audio_tx │ │
│ └──────────────────────┘ │
└─────────────────────────────────────────────────────────────────┘
数据面(RTP 转发):
Publisher ── RTP ──▶ TrackRemote ──▶ broadcast::Sender ──▶ TrackLocalStaticRTP ──▶ Subscriber(s)
关键设计决策:
- HTTP 而非 WebSocket 做 SDP 交换:每次 publish/subscribe 只需一个 Request/Response,无需维护额外长连接。
- 广播总线(
broadcast::channel):一个 Publisher 的 RTP 流可被任意数量的 Subscriber 消费,零拷贝广播。 - 信令与媒体共享
RoomRegistry:信令层(WebSocket)和媒体层(SFU)操作同一份房间状态,保证一致性。
3. 核心数据结构设计
3.1 PublisherBus:媒体广播总线
const BROADCAST_CAP: usize = 512;
#[derive(Clone)]
pub struct PublisherBus {
pub publisher_id: String,
pub video_tx: broadcast::Sender<RtpPacket>, // 视频广播发送端
pub audio_tx: broadcast::Sender<RtpPacket>, // 音频广播发送端
pub pli_tx: mpsc::Sender<()>, // PLI 关键帧请求通道
}
为什么用 broadcast::channel 而不是 Vec<mpsc::Sender>?
broadcast 的语义是“一写多读,每个读者独立消费”,完全匹配 SFU “一路上行 → 多路下行”的转发模型:
Publisher ──write──▶ broadcast::Sender
│
┌────────────┼────────────┐
▼ ▼ ▼
Receiver-1 Receiver-2 Receiver-3
(Subscriber) (Subscriber) (Subscriber)
新 Subscriber 只需调用 video_tx.subscribe() 即可接入,无需修改 Publisher 侧任何代码。
BROADCAST_CAP = 512 是缓冲区大小。超出时最老的包会被丢弃,Receiver 下次 recv() 会得到 RecvError::Lagged(n),此时触发 PLI 请求关键帧(详见第 8 节)。
3.2 Room:合并的房间状态
pub struct Room {
pub room_id: String,
/// 信令成员(WebSocket 连接)
pub signal_members: HashMap<String, SignalMember>,
/// 媒体 Publisher(推流者)
pub publishers: HashMap<String, PublisherBus>,
/// MLS KeyPackage 缓存
pub key_packages: HashMap<String, String>,
/// 订阅者的 PeerConnection(重协商用)
pub subscribers: HashMap<String, Arc<RTCPeerConnection>>,
/// 房间创建时间(计算通话时长)
pub started_at: chrono::DateTime<chrono::Utc>,
}
信令与媒体为何共用一个 Room?
分开管理的问题:信令系统需要知道“谁在推流”来决定是否广播 MemberJoined,媒体层需要知道“谁订阅了”来做重协商 Offer 推送。两套状态之间的同步是一个经典的分布式一致性问题,在单进程内完全没有必要引入这个复杂度。
合并的好处:
- 信令层和媒体层通过同一把
RwLock<HashMap<String, Room>>读写,天然一致 is_empty()可以同时检查信令成员、Publisher、Subscriber,精确判断房间是否可回收
3.3 AppState:依赖注入容器
#[derive(Clone)]
pub struct AppState {
pub rooms: RoomRegistry,
pub ice_servers: Vec<RTCIceServer>,
pub pool: sqlx::PgPool,
pub auth: AuthState,
pub claim_rate_limiter: Arc<RateLimiter>,
pub config: Arc<RwLock<RuntimeConfigRow>>,
pub sfu_metrics: Arc<metrics::SfuMetrics>,
pub global_sessions: GlobalSessionRegistry,
pub storage: Arc<dyn FileStorage>,
pub webrtc_config: WebRtcNetConfig, // 网络配置(公网IP/端口范围)
}
Axum 的 State<AppState> 通过 Clone 派发到每个 Handler,所有共享资源均通过 Arc 保证线程安全,无需手动管理生命周期。
4. PeerConnection 构建:网络配置详解
async fn build_pc(
ice_servers: Vec<RTCIceServer>,
net_config: &WebRtcNetConfig,
) -> Result<Arc<RTCPeerConnection>>
这是整个 SFU 的基础能力,每次 publish/subscribe 都会调用它创建一个新的 PeerConnection。
4.1 媒体编解码器注册
// VP8 视频,payload type 96
me.register_codec(RTCRtpCodecParameters {
capability: RTCRtpCodecCapability {
mime_type: MIME_TYPE_VP8.to_owned(),
clock_rate: 90000,
..
},
payload_type: 96,
..
}, RTPCodecType::Video)?;
// Opus 音频,payload type 111
me.register_codec(RTCRtpCodecParameters {
capability: RTCRtpCodecCapability {
mime_type: MIME_TYPE_OPUS.to_owned(),
clock_rate: 48000,
channels: 2,
sdp_fmtp_line: "minptime=10;useinbandfec=1".to_owned(),
..
},
payload_type: 111,
..
}, RTPCodecType::Audio)?;
useinbandfec=1:启用 Opus 带内前向纠错,弱网下音质显著提升。
minptime=10:最小打包间隔 10ms,平衡延迟与压缩效率。
4.2 云服务器 NAT 穿透配置
这是生产部署最容易踩坑的地方。云服务器(阿里云/腾讯云/AWS)普遍使用 1:1 EIP NAT,容器内部看到的是私网 IP,但 ICE Candidate 必须携带公网 IP 才能被客户端连通。
// 只允许 IPv4 UDP,避免 IPv6 ICE 收集干扰
se.set_network_types(vec![NetworkType::Udp4]);
// 告知 webrtc-rs 真实公网地址
if let Some(ref public_ip) = net_config.public_ip {
se.set_nat_1to1_ips(
vec![public_ip.clone()],
RTCIceCandidateType::Host, // 直接声明为 Host Candidate
);
}
// 限制 UDP 端口范围,让安全组精确放行
if let Some((min, max)) = net_config.udp_port_range {
let ephemeral = EphemeralUDP::new(min, max)?;
se.set_udp_network(UDPNetwork::Ephemeral(ephemeral));
}
配置了公网 IP 后,服务端不再需要 STUN:
let server_ice_servers = if net_config.public_ip.is_some() {
vec![] // 公网 IP 已知,不需要查 STUN
} else {
ice_servers // 本地开发环境,使用配置的 STUN 服务器
};
这样可以避免服务端去访问墙外 STUN 造成的额外延迟。
对应的 production.yaml 配置:
webrtc:
public_ip: "1.2.3.4" # 服务器公网 IP
udp_port_min: 40000
udp_port_max: 40100
安全组放行规则(以阿里云为例):
| 方向 | 协议 | 端口范围 | 说明 |
|---|---|---|---|
| 入方向 | UDP | 40000-40100 | WebRTC 媒体流 |
| 入方向 | TCP | 443 | HTTPS/WSS 信令 |
4.3 ICE 收集超时
tokio::time::timeout(
std::time::Duration::from_secs(timeout_secs),
rx.recv()
).await.map_err(|_| anyhow::anyhow!(
"ICE Candidate 收集超时({}s),请检查:\n\
1. production.yaml 中 webrtc.public_ip 是否填写正确\n\
2. 安全组 UDP 端口 {}-{} 是否已开放\n\
3. STUN 服务器是否可达",
...
))?;
超时时直接报错并给出具体排查方向,比返回一个模糊的 500 错误对运维友好得多。
5. Publish 流程
Client Server
│ │
│── POST /room/{id}/publish/{uid} ──▶│
│ body: base64(SDP Offer) │
│ │ 1. JWT 鉴权
│ │ 2. build_pc()
│ │ 3. add_transceiver(Video, Recvonly)
│ │ 4. add_transceiver(Audio, Recvonly)
│ │ 5. room.add_publisher() → PublisherBus
│ │ 6. pc.on_track() 注册回调
│ │ 7. sdp_exchange() 完成协商
│◀── 200 OK: base64(SDP Answer) ────│
│ │
│══════════ UDP RTP 媒体流 ══════════│
│ │ on_track 触发:
│ │ tokio::spawn(run_publisher_track)
5.1 Transceiver 方向设置
pc.add_transceiver_from_kind(
RTPCodecType::Video,
Some(RTCRtpTransceiverInit {
direction: RTCRtpTransceiverDirection::Recvonly, // 服务端只收
send_encodings: vec![],
}),
).await?;
Recvonly 是关键:服务端明确声明只接收,不发送。SDP 协商后浏览器/客户端就知道该往这条 PeerConnection 推流。
5.2 on_track 回调注册
pc.on_track(Box::new(
move |track: Arc<TrackRemote>, _receiver, _transceiver| {
let video_tx = video_tx.clone();
let audio_tx = audio_tx.clone();
Box::pin(async move {
// 视频轨道才需要 PLI 接收器
let pli_rx = if track.kind() == RTPCodecType::Video {
pli_rx_ref.lock().await.take().unwrap_or_else(|| {
let (_, rx) = mpsc::channel(1);
rx
})
} else {
let (_, rx) = mpsc::channel(1);
rx
};
tokio::spawn(peer::run_publisher_track(
track, video_tx, audio_tx, pc_pli, pli_rx, metrics,
));
})
},
));
注意 pli_rx_ref.lock().await.take():PLI 接收器只应该被视频轨道取走一次(Option::take),音频轨道拿到的是一个占位的永不触发的 channel,避免资源泄漏。
6. Subscribe 流程
Client Server
│ │
│── POST /room/{id}/subscribe/{uid}──▶│
│ body: base64(SDP Offer) │
│ │ 1. JWT 鉴权
│ │ 2. 获取其他所有 Publisher 快照
│ │ 3. build_pc()
│ │ 4. 为每个 Publisher:
│ │ a. 创建 TrackLocalStaticRTP
│ │ b. pc.add_track()
│ │ c. spawn(run_subscriber_video/audio)
│ │ d. bus.request_keyframe()
│ │ 5. room.subscribers.insert(pc)
│ │ 6. sdp_exchange()
│◀── 200 OK: base64(SDP Answer) ────│
│ │
│◀═════════ UDP RTP 媒体流 ══════════│
6.1 “幽灵订阅”设计
let publishers = {
let registry = state.rooms.read().await;
match registry.get(&room_id) {
Some(room) => room.other_publishers(&user_id),
None => vec![], // 房间不存在也返回空,允许建立幽灵订阅连接
}
};
允许在没有任何 Publisher 的情况下建立 Subscribe 连接,原因是:
用户加入顺序不确定。在实际使用中,Alice 可能比 Bob 慢几百毫秒才开始推流,如果订阅接口直接返回错误,客户端就需要实现重试逻辑,增加复杂度。
“幽灵订阅”方案:先建立 PC 并存入 room.subscribers,当 Bob 后来推流时,SFU 在 handle_publish 中检测到已有订阅者,主动发起重协商(见第 7 节)。
6.2 Track 到 Stream 的映射
let video_track = Arc::new(TrackLocalStaticRTP::new(
RTCRtpCodecCapability {
mime_type: MIME_TYPE_VP8.to_owned(), ..
},
format!("video-{}", video_ssrc), // track_id:全局唯一
format!("stream-{}", bus.publisher_id), // stream_id:标识媒体流归属
));
stream_id 格式为 "stream-{publisher_id}",前端 JavaScript 通过 track.streams[0].id 解析出 publisher_id,从而将视频轨道与对应用户的 UI 元素绑定:
// frontend/src/webrtc.ts
pc.ontrack = (event) => {
const streamId = event.streams[0]?.id; // "stream-alice"
const publisherId = streamId?.replace("stream-", ""); // "alice"
bindVideoElement(publisherId, event.track);
};
7. 动态新成员加入:重协商机制
这是 SFU 中最复杂的场景:Alice 已经在房间里订阅,Bob 此时才加入并开始推流,如何让 Alice 收到 Bob 的媒体流?
7.1 流程图
Bob 开始推流
│
▼
handle_publish() 执行成功
│
├── tokio::spawn(async move { ← 异步,不阻塞 publish 响应
│ sleep(1s) ← 等 Bob 自己先建联
│
│ 1. 获取 subscribers 快照(短暂持读锁)
│ 2. 释放全局读锁 ←── 关键!锁外执行耗时协商
│
│ for each subscriber (Alice):
│ a. 创建新的 TrackLocalStaticRTP
│ b. sub_pc.add_track(video_track)
│ c. sub_pc.add_track(audio_track)
│ d. spawn(run_subscriber_video/audio) ← 开始转发
│ e. sub_pc.create_offer()
│ f. sub_pc.set_local_description()
│ g. room.send_to(alice, Offer{sdp}) ← 通过 WS 推送
│ })
│
▼
Alice 收到 WS Offer
│
▼
Alice 前端: pc.setRemoteDescription(offer) → createAnswer() → sendAnswer()
│
▼
POST /room/{id}/subscribe/{alice}/answer
│
▼
handle_subscribe_answer():
sub_pc.set_remote_description(answer) ← 重协商完成
7.2 锁外执行的重要性
// ✅ 正确:先拿快照,立刻释放锁,锁外做耗时操作
let subscribers_snapshot: Vec<(String, Arc<RTCPeerConnection>)> = {
let registry = rooms_clone.read().await;
registry.get(&room_id_clone)
.map(|room| room.subscribers.iter()
.filter(|(id, _)| *id != &pub_user_id)
.map(|(id, pc)| (id.clone(), Arc::clone(pc)))
.collect())
.unwrap_or_default()
}; // ← 读锁在此释放!
// 锁释放后,慢慢执行耗时的 SDP 协商(可能 100ms-1s)
for (sub_id, sub_pc) in subscribers_snapshot {
// ... create_offer / set_local_description / send via WS
}
如果在持锁状态下执行 create_offer(),整个 RoomRegistry 的读锁会被持有数百毫秒,期间其他任何需要写锁的操作(如新用户加入)都会被阻塞。
7.3 Answer 回路
async fn handle_subscribe_answer(
Path((room_id, user_id)): Path<(String, String)>,
State(state): State<AppState>,
headers: HeaderMap,
body: String,
) -> HandlerResult<String> {
verify_sfu_request(&headers, &state, &user_id)?;
let answer = decode_sdp(&body)?;
let pc = {
let registry = state.rooms.read().await;
registry.get(&room_id)
.and_then(|room| room.subscribers.get(&user_id))
.cloned()
.ok_or_else(|| AppError(anyhow::anyhow!("Subscriber not found")))?
};
pc.set_remote_description(answer).await?;
Ok("ok".to_string())
}
这个接口非常简单:从 room.subscribers 取出 Alice 的 PC,应用她返回的 Answer,重协商完成。
8. RTP 转发与 PLI 自愈
8.1 Publisher 轨道处理
pub async fn run_publisher_track(
track: Arc<TrackRemote>,
video_tx: broadcast::Sender<RtpPacket>,
audio_tx: broadcast::Sender<RtpPacket>,
pc: Arc<RTCPeerConnection>,
mut pli_rx: mpsc::Receiver<()>,
metrics: Arc<SfuMetrics>,
) {
// ...
let mut buf = vec![0u8; 1500]; // MTU 标准大小
loop {
match track.read(&mut buf).await {
Ok((packet, _attr)) => {
let approx_size = 12 + packet.payload.len() as u64;
metrics.add_recv(approx_size);
let tx = if kind == RTPCodecType::Video { &video_tx } else { &audio_tx };
let _ = tx.send(packet); // 忽略错误:没有订阅者时 send 返回 Err 是正常的
}
Err(e) => {
info!("Publisher track ended: {}", e);
break;
}
}
}
// 函数退出时 pli_rx drop,PLI 协程自动退出(None 分支)
}
设计亮点:
pli_rx在run_publisher_track退出时自动 drop,导致 PLI 协程的pli_rx.recv()返回None,优雅地级联关闭相关协程,无需显式取消信号。
8.2 PLI 双触发机制
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(3));
loop {
tokio::select! {
// 触发器 1:定时 PLI(每 3 秒),保证新订阅者定期收到完整帧
_ = interval.tick() => {
send_pli(&pc_pli, ssrc_pli).await;
}
// 触发器 2:按需 PLI(新订阅者加入时立刻请求)
msg = pli_rx.recv() => {
match msg {
Some(()) => {
info!("PLI triggered by new subscriber");
send_pli(&pc_pli, ssrc_pli).await;
}
None => break, // channel 关闭,退出
}
}
}
}
});
PLI(Picture Loss Indication):RTCP 反馈消息,告知编码器“我丢包了,请发一个完整的 I 帧”。
两种触发场景:
- 定时触发(3s):确保新加入的订阅者在最多 3 秒内看到完整画面
- 按需触发:新订阅者加入时立刻发送,让他更快看到清晰画面;
broadcast::Lagged时触发,自愈花屏问题
8.3 Subscriber 视频转发与自愈
pub async fn run_subscriber_video(
track: Arc<TrackLocalStaticRTP>,
mut rx: broadcast::Receiver<RtpPacket>,
viewer_ssrc: u32,
metrics: Arc<SfuMetrics>,
pli_tx: mpsc::Sender<()>,
) {
loop {
match rx.recv().await {
Ok(mut packet) => {
packet.header.ssrc = viewer_ssrc; // 重写 SSRC,避免冲突
if let Err(e) = track.write_rtp(&packet).await {
warn!("write error: {}", e);
break;
}
metrics.add_send(12 + packet.payload.len() as u64);
}
Err(broadcast::error::RecvError::Lagged(n)) => {
warn!("Subscriber video lagged {} packets", n);
// 丢包 → 花屏 → 请求关键帧自愈
let _ = pli_tx.try_send(());
}
Err(broadcast::error::RecvError::Closed) => break,
}
}
}
SSRC 重写:每个 Subscriber 得到一个独立的随机 SSRC,避免多个订阅者之间的 SSRC 冲突导致浏览器混淆媒体流。
Lagged 自愈链路:
广播缓冲区满(BROADCAST_CAP=512)
│
▼
Subscriber recv() 返回 Lagged(n)
│
▼
try_send(()) → pli_tx
│
▼
run_publisher_track 中的 pli_rx.recv()
│
▼
send_pli(pc, ssrc) → Publisher 端发送 I 帧
│
▼
Subscriber 收到完整 I 帧,画面恢复
9. 带宽监控:无锁原子计数器
9.1 数据结构
#[derive(Debug)]
pub struct SfuMetrics {
pub recv_bytes: AtomicU64, // Publisher 上行总字节
pub send_bytes: AtomicU64, // Subscriber 下行总字节
current_mbps_bits: AtomicU64, // 当前带宽(f64 bits 形式存储)
}
为什么用 Relaxed ordering?
metrics.add_recv(approx_size);
// ↓ 实现
self.recv_bytes.fetch_add(bytes, Ordering::Relaxed);
带宽统计是一个 “最终一致” 场景,我们不需要跨线程同步保证(不需要 Acquire/Release),只需要原子性(不丢数据)。Relaxed 在 x86 上几乎等于普通加法,没有内存屏障开销。
f64 原子存储技巧:
// Rust 标准库没有 AtomicF64,通过 bits 转换解决
metrics_clone.current_mbps_bits
.store(mbps.to_bits(), Ordering::Relaxed);
// 读取时
f64::from_bits(self.current_mbps_bits.load(Ordering::Relaxed))
9.2 后台采样任务
pub fn start_sampling_loop(self: &Arc<Self>, interval_secs: u64) {
let metrics_clone = Arc::clone(self);
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(interval_secs));
interval.set_missed_tick_behavior(MissedTickBehavior::Skip); // 跳过积压的 tick
let mut last_recv = 0u64;
let mut last_send = 0u64;
let mut last_time = Instant::now();
loop {
interval.tick().await;
let now = Instant::now();
let elapsed = now.duration_since(last_time).as_secs_f64();
let curr_recv = metrics_clone.recv_bytes.load(Ordering::Relaxed);
let curr_send = metrics_clone.send_bytes.load(Ordering::Relaxed);
let mbps = if elapsed > 0.1 {
// (字节差 * 8位/字节) / 1_000_000 / 秒 = Mbps
((curr_recv - last_recv + curr_send - last_send) as f64 * 8.0)
/ 1_000_000.0 / elapsed
} else {
0.0
};
metrics_clone.current_mbps_bits.store(mbps.to_bits(), Ordering::Relaxed);
last_recv = curr_recv;
last_send = curr_send;
last_time = now;
}
});
}
MissedTickBehavior::Skip:如果采样任务本身因为系统繁忙而延迟,跳过积压的 tick 而不是突发补偿,避免“惊群”效应。
10. 鉴权设计:防越权攻击
fn verify_sfu_request(
headers: &HeaderMap,
state: &AppState,
path_user_id: &str,
) -> anyhow::Result<String> {
// 1. 提取 Bearer Token
let token = extract_bearer_token(headers)
.ok_or_else(|| anyhow::anyhow!("SFU 请求缺少 Authorization Token (鉴权失败)"))?;
// 2. 验证 JWT 签名和有效期
let claims = state.auth.jwt.verify(&token)
.map_err(|e| anyhow::anyhow!("Token 无效 (鉴权失败): {}", e))?;
// 3. 关键:路径 user_id 必须与 JWT sub 一致
if claims.sub != path_user_id {
anyhow::bail!(
"越权操作:Token 归属 {},但尝试操作 {} 的媒体流",
claims.sub,
path_user_id
);
}
Ok(claims.sub)
}
为什么需要第 3 步?
仅验证 JWT 有效性是不够的。攻击场景:
Bob 持有自己的有效 JWT
Bob 发送 POST /room/xxx/publish/alice
↓
如果只检查 JWT 有效性:Bob 可以以 alice 身份推流!
↓
正确做法:JWT.sub 必须等于路径里的 user_id
错误码映射:
impl IntoResponse for AppError {
fn into_response(self) -> Response {
let status = if self.0.to_string().contains("鉴权失败") {
StatusCode::UNAUTHORIZED // 401
} else if self.0.to_string().contains("越权操作") {
StatusCode::FORBIDDEN // 403
} else {
StatusCode::INTERNAL_SERVER_ERROR // 500
};
(status, format!("SFU error: {}", self.0)).into_response()
}
}
通过错误信息中的关键词区分 401/403/500,避免将所有错误都返回 500 给客户端。
11. 房间生命周期管理
11.1 状态变化回调
Publisher 和 Subscriber 都注册了 on_peer_connection_state_change:
pc.on_peer_connection_state_change(Box::new(move |s| {
Box::pin(async move {
match s {
// ⚠️ 注意:Disconnected 不能清理!
// Disconnected 表示 ICE 暂时失联(网络抖动),可能自动恢复
// 只有 Failed / Closed 才是真正结束
RTCPeerConnectionState::Failed | RTCPeerConnectionState::Closed => {
let mut registry = rooms_for_cleanup.write().await;
if let Some(room) = registry.get_mut(&room_id) {
room.remove_publisher(&user_id); // or subscribers.remove()
if room.is_empty() {
registry.remove(&room_id);
info!("Room [{}] cleaned up", room_id);
}
}
}
_ => {}
}
})
}));
Disconnected 陷阱:很多初学者在看到 Disconnected 时立刻清理资源,导致用户手机稍微切换一下网络(从 WiFi 到 4G)就被踢出通话。正确做法是只在 Failed 或 Closed 时清理。
11.2 is_empty() 的精确判断
pub fn is_empty(&self) -> bool {
self.signal_members.is_empty()
&& self.publishers.is_empty()
&& self.subscribers.is_empty()
}
三个集合都为空才回收房间,避免以下场景的误回收:
- 有人在 WS 连接但还没推流(
signal_members不空) - 有人订阅但推流者还没加入(
subscribers不空)
12. 管控能力:强踢与会话终止
12.1 强制下线用户
面向管理后台(gRPC 接口),可强制将用户踢下线:
pub async fn force_disconnect_user(
global_sessions: &GlobalSessionRegistry,
room_registry: &RoomRegistry,
user_id: &str,
reason: &str,
) {
// 步骤一:发 SystemKick(必须在清理 signal_members 之前!)
// 因为清理后就找不到这个人的 WS 连接了
{
let sessions = global_sessions.write().await;
if let Some(session) = sessions.get(user_id) {
let _ = session.tx.send(ServerMessage::SystemKick {
reason: reason.to_string(),
});
}
} // 写锁立即释放
// 步骤二、三:清理房间状态,收集需要关闭的 PC
let mut pcs_to_close = Vec::new();
{
let mut registry_lock = room_registry.write().await;
for (room_id, room) in registry_lock.iter_mut() {
// ... 清理 signal_members / publishers / subscribers
if let Some(pc) = room.subscribers.remove(user_id) {
pcs_to_close.push(pc); // 收集,不在锁内 close
}
}
// 回收空房间
registry_lock.retain(|_, room| !room.is_empty());
} // 写锁释放
// 步骤四:锁外关闭 WebRTC 连接(异步 IO,不能持锁)
for pc in pcs_to_close {
let _ = pc.close().await;
}
}
为什么不能在锁内 close() PC?
pc.close() 是异步操作,内部会等待 ICE 和 DTLS 的关闭握手(可能需要几百毫秒)。如果在持有 room_registry 写锁期间调用,整个房间状态对所有其他协程不可用,严重影响并发性能。
正确做法:持锁期间只做 remove(),将 PC 收集到临时 Vec,锁释放后再逐一 close()。
12.2 终止指定房间的通话
pub async fn terminate_session(
registry: &RoomRegistry,
room_id: &str,
reason: &str,
) -> bool {
let mut pcs_to_close = Vec::new();
{
let mut registry_lock = registry.write().await;
let room = match registry_lock.get_mut(room_id) {
Some(r) => r,
None => return false,
};
// 1. 通知所有人(保留 WS 连接,只终止通话)
for member in room.signal_members.values() {
let _ = member.tx.send(ServerMessage::SessionTerminated {
room_id: room_id.to_string(),
reason: reason.to_string(),
});
}
// 2. 清空媒体资源(drop PublisherBus → broadcast channel 自动关闭)
room.publishers.clear();
// 3. 收集 subscriber PC
for pc in room.subscribers.values() {
pcs_to_close.push(Arc::clone(pc));
}
room.subscribers.clear();
}
// 4. 锁外关闭
for pc in pcs_to_close {
let _ = pc.close().await;
}
true
}
与 force_disconnect_user 的区别:
force_disconnect_user:踢人,断 WS 连接,用户感知到被管理员踢出系统terminate_session:只终止通话,WS 连接保留,用户感知到“通话被管理员终止”,可以继续聊天
13. 生产环境部署要点
13.1 最小化配置清单
# configuration/production.yaml
webrtc:
public_ip: "YOUR_SERVER_PUBLIC_IP"
udp_port_min: 40000
udp_port_max: 40100
13.2 安全组 / 防火墙规则
# UFW 示例
ufw allow 443/tcp # HTTPS + WSS
ufw allow 40000:40100/udp # WebRTC 媒体
# iptables 示例
iptables -A INPUT -p udp --dport 40000:40100 -j ACCEPT
13.3 性能基准参考
在 4 核 8G 云服务器上(无 SFU 硬件加速),单房间性能参考:
| 房间人数 | 上行带宽(VP8 360p) | 下行带宽(合计) | CPU 占用 |
|---|---|---|---|
| 2 人 | ~400 Kbps | ~400 Kbps | < 5% |
| 4 人 | ~1.6 Mbps | ~4.8 Mbps(每人收3路) | ~15% |
| 8 人 | ~3.2 Mbps | ~22 Mbps(每人收7路) | ~40% |
SFU 的下行带宽 = N × (N-1) × 单路上行带宽,人数越多下行压力越大,这是 SFU 相比 MCU 的主要劣势。
13.4 关键超时参数
// ICE 收集超时(sdp_exchange 函数中)
let timeout_secs = 10u64;
// 新成员推流后等待建联的时间(handle_publish 重协商任务中)
tokio::time::sleep(Duration::from_secs(1)).await;
// PLI 定时发送间隔(run_publisher_track 中)
let mut interval = tokio::time::interval(Duration::from_secs(3));
// 广播缓冲区大小(room.rs 中)
const BROADCAST_CAP: usize = 512;
14. 已知局限与后续规划
14.1 当前局限
| 问题 | 说明 | 影响 |
|---|---|---|
| 单进程状态 | RoomRegistry 存内存,多实例不共享 | 无法水平扩展 |
| 无 SVC/Simulcast | 只转发单一质量流 | 弱网用户体验差 |
| 无带宽自适应(BWE) | 不根据网络动态调整码率 | 弱网可能卡顿 |
| VP8 Only | 未支持 H.264/AV1 | 部分设备硬件解码无法利用 |
| TCP 回退缺失 | UDP 不通时无 TURN over TCP | 严格防火墙下无法连接 |
14.2 后续规划
近期(可独立实现):
- 引入
TURN服务器(coturn),解决严格防火墙场景 Simulcast支持:Publisher 推多路质量,SFU 按需选择转发
中期(需要架构改造):
- Redis 发布订阅:将
broadcast::channel升级为跨进程广播,支持多节点 SFU - 基于 REMB/TWCC 的带宽估算,实现服务端 BWE
- Prometheus 指标导出(当前
SfuMetrics已有原子计数器,接入成本低)
长期(较大工程量):
- SFU 级联(大型会议室多节点转发)
小结
本文详细介绍了基于 webrtc-rs + Axum + Tokio 构建 SFU 服务的核心实现,关键设计要点回顾:
| 模块 | 核心设计 |
|---|---|
| 媒体广播 | broadcast::channel 一写多读,SSRC 重写避免冲突 |
| 动态加入 | “幽灵订阅” + 锁外重协商,避免全局锁死锁 |
| PLI 自愈 | 定时 + 按需双触发,Lagged 时自动请求关键帧 |
| NAT 穿透 | set_nat_1to1_ips + 端口范围限制,适配云服务器 EIP |
| 鉴权安全 | JWT + path_user_id 双重校验,防止越权推流 |
| 管控能力 | 分步清理(发信令→清状态→锁外关PC),避免锁内 IO |
| 带宽监控 | AtomicU64 + Relaxed ordering,零锁开销统计 |
代码仓库:(链接)
如有问题欢迎在评论区交流,也欢迎 Star/Fork 项目参与共建。
系列导航
火山引擎视频云技术社区,是面向 AI 音视频开发者的技术交流平台。这里汇聚源自抖音、豆包等亿级 DAU 产品的 RTC、直播、点播、AI 媒体处理、音视频互动技术,提供接入指南、最佳实践、性能调优、场景案例、Demo 代码、开源项目、白皮书和 API 文档。社区汇聚官方工程师与一线开发者,为 AI 视频通话、数字人、AI 视频处理等应用的开发与落地提供技术支持。
更多推荐
所有评论(0)