Вход на сайт

Просмотр новости

Найдите то, что Вас интересует

Заметки о том, как я писал SFU на Rust (2 часть)

Дата публикации: 12-08-2026 13:42:30

Продолжаем изучать, как устроен SFU изнутри: добавляем simulcast, ice restart и мониторинг сети. Дневник разработки на Rust. Читать далее

Основное содержимое страницы с новостью.

Большое спасибо всем, кто прочитал первую часть, в этой части мы продолжим разбираться с webrtc и реализуем действительно полезные возможности SFU.

Введение

В прошлой части мы реализовали сигналинг и базовую маршрутизацию видео трансляций между пользователями. В этой части мы реализуем:

  • Мониторинг сети - отслеживание статистики подключений и сохранение ее для аналитики.

  • Simulcast - одновременная отправка видео в нескольких качествах, чтобы SFU мог выбирать подходящий поток для каждого зрителя в зависимости от его канала.

  • Ice restart - алгоритм восстановления соединения.

Но перед тем как перейти к реализации новых фичей нашего приложения, нам нужно убедиться, что по итогу у нас не сломается старый функционал. Поэтому пора написать немного тестов.

Тестирование

Чтобы протестировать приложение, при этом не переоткрывая его по сто раз в браузере, напишем тестовый клиент, который будет симулировать работу нашего фронтенда.

pub struct TestClient {
    pub peer_id: Option<Uuid>,
    pub publisher_pc: Arc<RTCPeerConnection>,
    pub subscriber_pc: Arc<RTCPeerConnection>,
    pub sfu_tx: tokio::sync::mpsc::Sender<SignalMessage>,
    pub sfu_rx: tokio::sync::mpsc::Receiver<SignalMessage>,
    pub connected_peers: Arc<Mutex<HashSet<Uuid>>>,
    pub markers: Markers,
}

pub struct Markers {
    pub ice_connected: bool,
    pub send_offer: bool,
    pub receive_answer: bool,
    pub receive_offer: bool,
}

В тестах мы не хотим общаться между клиентом и сервером посредством вебсокетов, поэтому введем трейт, который послужит абстракцией для коммуникации.

pub trait SyncChannel: Send + 'static {
    type Item: From<SignalMessage>;
    type Error: std::error::Error + Debug;
    fn send(&mut self, message: Self::Item) -> 
    impl Future<Output = Result<(), Self::Error>> + Send;
}

И теперь заменим конкретную реализацию на generic-параметр у нашего пользователя.

pub struct User<S: SyncChannel> {
    pub room: Addr<Room<S>>,
    pub peer_id: Uuid,
    // before pub sync_channel: Websocket
    pub sync_channel: S,
    pub publisher: WeakAddr<Publisher<S>>,
    pub subscriber: WeakAddr<Subscriber<S>>,
}

Отныне в тестах мы можем использовать обычные mpsc каналы для сигналинга.

pub struct TestSyncChannel {
    channel: tokio::sync::mpsc::Sender<SignalMessage>                               
}

impl SyncChannel for TestSyncChannel {
    type Item = SignalMessage;
    type Error = tokio::sync::mpsc::error::SendError<SignalMessage>;
    async fn send(&mut self, message: SignalMessage) -> Result<(), Self::Error> {
        self.channel.send(message)
            .await?;
        Ok(())
    }
}

Вернемся к нашему тестовому клиенту, в отличие от браузера мы будем стримить обычный видео-файл формата .ivf. Поскольку чтение и декодирование это синхронные операции, а запись данных в трек асинхронная, разделим их на tokio::task::spawn_blocking и tokio::task::spawn.

tokio::task::spawn_blocking(move || {
    let file = File::open("../output.ivf")?;
    let buf_reader = BufReader::new(file);
    let (mut reader, header) = webrtc::media::IVFReader::new(buf_reader)?;
    let sleep_time = Duration::from_millis(
        ((1000 * header.timebase_numerator) / header.timebase_denominator) as u64,
    );
    loop {
        let frame = match reader.parse_next_frame() {
            Ok((frame, _)) => frame,
            Err(_) => {
                let file = File::open("../output.ivf")?;
                let buf_reader = BufReader::new(file);
                reader = IVFReader::new(buf_reader)?.0;
                continue;
            }
        };
        if tx.blocking_send((frame, sleep_time)).is_err() {
            break ;
        }
        std::thread::sleep(sleep_time);
    }
    Ok(())
});
tokio::task::spawn(async move {
    while let Some((frame, duration)) = rx.recv().await {
        video_track
            .sample_writer()
            .write_sample(&Sample {
                data: frame.freeze(),
                duration,
                ..Default::default()
            })
            .await?;
    }
    Result::<_, Error>::Ok(())
});

Реализуем обработку сигналинга на нашем тестовом клиенте.

match message {
  SignalMessage::Welcome { peer_id } => self.send_track(peer_id).await?,
  SignalMessage::PeerLeft { peer_id }  => {
      let mut peers = self.connected_peers.lock().await;
      peers.remove(&peer_id);
  }
  SignalMessage::Rtc { target, message_type } => {
      let pc = match target {
          sfu::user::Target::Publisher => &self.publisher_pc,
          sfu::user::Target::Subscriber => &self.subscriber_pc
      };
      match message_type {
          sfu::user::MessageType::Candidate { candidate } => {
              self.handle_ice_candidate(pc, candidate).await?;
              self.markers.ice_connected = true;
          }
          sfu::user::MessageType::Answer { sdp } => {
              let description = RTCSessionDescription::answer(sdp)?;
              self.publisher_pc.set_remote_description(description).await?;
              self.markers.receive_answer = true;
          }
          sfu::user::MessageType::Offer { sdp } => {
              self.handle_offer(sdp).await?;
              self.markers.receive_offer = true;
          },
          sfu::user::MessageType::IceRestart { sdp } => 
            self.handle_ice_restart(target, pc, sdp).await?
      }
  },
}

Напишем тест для проверки подключения одного пользователя к серверу.

#[tokio::test(flavor = "multi_thread")]
async fn connect_one_user_to_room() -> Result<(), Error> {
    tracing_subscriber::fmt::fmt()
        .with_env_filter("webrtc=ERROR,webrtc_ice=ERROR")
        .init();
    let server: sfu::actor::Addr<Server<_>> = Server::default().start();
    let (tx, rx) = tokio::sync::oneshot::channel();
    let _ = server.send(ServerMessage::CreateRoom { name: "".to_string(), response_channel: tx }).await;
    let (_, room) = rx.await?;
    let (channel, mut client, stream) = spawn_test_client().await?;
    let user = User::new(channel, room.clone()).await?;
    let user_id = user.peer_id;
    let addr = user.start();
    addr.add_stream(tokio_stream::wrappers::ReceiverStream::new(stream));
    let _ = room.send(RoomMessage::Join { peer_id: user_id, addr }).await;
    client.setup().await?;
    for _stage in 0 .. 2 {
        client.handle_message().await?;
    } 
    assert_eq!(client.markers.send_offer, true);
    assert_eq!(client.markers.receive_answer, true);
    Ok(())
}
running 1 test
test connect_to_room::connect_one_user_to_room ... ok
test result: ok. 1 passed; 0 failed; 0 ignored; 0 measured; 3 filtered out; finished in 0.02s

А также для проверки подключения нескольких пользователей к серверу и к друг другу.

#[tokio::test]
async fn connect_many_users_to_room() -> Result<(), Error> {
    tracing_subscriber::fmt::fmt()
        .with_env_filter("webrtc=ERROR,webrtc_ice=ERROR")
        .init();
    let server: sfu::actor::Addr<Server<_, FileStorage>> = Server::default().start();
    let (tx, rx) = tokio::sync::oneshot::channel();
    let _ = server.send(ServerMessage::CreateRoom { name: "".to_string(), response_channel: tx }).await;
    let (_, room) = rx.await.unwrap();
    let mut tasks = Vec::with_capacity(24);
    let barrier = Arc::new(tokio::sync::Barrier::new(24));
    for _ in 0 .. 24 {
        let (channel, mut client, stream) = spawn_test_client().await?;
        let user = User::new(channel, room.clone()).await?;
        let peer_id = user.peer_id;
        let user = user.start();
        user.add_stream(tokio_stream::wrappers::ReceiverStream::new(stream));
        let _ = room.send(RoomMessage::Join { peer_id, addr: user }).await;
        let barrier = barrier.clone();        
        let task = tokio::spawn(async move {
            client.setup().await?;
            for _stage in 0 .. 4 {
                client.handle_message().await?;
            } 
            assert_eq!(client.markers.send_offer, true);
            assert_eq!(client.markers.ice_connected, true);
            assert_eq!(client.markers.receive_answer, true);
            assert_eq!(client.markers.receive_offer, true);
            barrier.wait().await;
            Result::<_, Error>::Ok(())
        });
        tasks.push(task);
    }
    futures_util::future::try_join_all(tasks).await.unwrap();
    Ok(())
}
running 1 test
test connect_to_room::connect_many_users_to_room ... ok
test result: ok. 1 passed; 0 failed; 0 ignored; 0 measured; 3 filtered out; finished in 2.50s

Тесты показывают, что наш функционал (пока) работает, начнем добавлять изменения!

Мониторинг сети

Мониторинг понадобится нам как для аналитики того, как работает наше приложение, так и для симулькаста, но о нем позже.

Для начала введем абстракцию для хранения записей о статистике.

pub trait Storage: Send + 'static + Sized {
    type Configuration: StorageConfiguration;
    type Error: std::error::Error + Debug;
    fn connect(configuration: &Self::Configuration) 
      -> impl Future<Output = Result<Self, Self::Error>> + Send;
    fn insert(&mut self, item: StorageItem) 
      -> impl Future<Output = Result<(), Self::Error>> + Send;
}

pub trait StorageConfiguration: Send +'static + Sized {
    type Error: std::error::Error + Debug;
    fn from_env() -> Result<Self, Self::Error>;
}

В реальном приложении я думаю использовал бы что то вроде Сlickhouse или же вообще подумал о Kafka и самими данными управлял бы в отдельном микросервисе, здесь же используем просто запись в файл.

pub struct FileStorage {
    file: tokio::fs::File
}

pub struct FileConfiguration { path: String }
Актор QualityMonitor

Перейдем к самому актору, который будет отслеживать статистику соединения пользователя.

pub struct QualityMonitor<C: SyncChannel, S: Storage> {
    publisher_connection: Arc<RTCPeerConnection>,
    id: Uuid,
    user: Addr<User<C, S>>,
    storage: S,
    last_stats_time: Instant,
    consecutive_high_signals: usize, 
    current_quality: Option<StreamQuality>, 
    current_stats: CurrentStats,
    thresholds: QualityThresholds,
    update_period: Duration,
}

При запуске актора, запустим стрим, который в интервале будет отправлять сигнал, о том что нужно собрать статистику.

async fn starting(&mut self, ctx: &crate::actor::Ctx<'_, Self>) {
    let stream = IntervalStream::new(interval(self.update_period));
    ctx.addr.add_stream(stream, |_| StreamItem::Next(QualityMonitorMessage::Ping));
}

В начале запросим у соединения статистику и посчитаем нужные метрики.

async fn update_stats(&mut self) {
    let stats = self.publisher_connection.get_stats().await;
    let now = std::time::Instant::now();
    let elapsed_secs = now.duration_since(self.last_stats_time).as_secs_f64();
    self.last_stats_time = now;
    // 1. Собираем метрики из отчетов WebRTC
    let (total_bytes, total_packets, total_nacks) = self.collect_video_metrics(&stats);
    // 2. Рассчитываем битрейт
    self.calculate_bitrate(total_bytes, elapsed_secs);
    // 3. Рассчитываем потери пакетов
    self.calculate_packet_loss(total_packets, total_nacks);
    // 4. Сохраняем состояние для следующего шага
    self.current_stats.last_bytes_received = total_bytes;
    self.current_stats.last_packets_received = total_packets;
    self.current_stats.last_nack_count = total_nacks;
}

Сохраним:

async fn save_stats(&mut self) -> Result<(), Error> {
    let item = crate::StorageItem { 
        connection_id: self.id, 
        stats: &self.current_stats, 
        timestamp: Utc::now(),  
    };
    self.storage.insert(item)
        .await
        .map_err(|e| Error::SystemError { message: e.to_string().into() })?;
    Ok(())
}
77ae00274ec691c0c9b70ed38a362779.png

После обновление метрик, определим какое качество соединения у пользователя.

fn get_quality(&self) -> Option<StreamQuality> {
    if self.current_stats.last_bytes_received == 0 {
        return None
    }
    if self.current_stats.packet_loss > self.thresholds.low_loss || 
       self.current_stats.bitrate_bps < self.thresholds.low_bitrate {
        return Some(StreamQuality::Low);
    }
    if self.current_stats.packet_loss > self.thresholds.mid_loss || 
       self.current_stats.bitrate_bps < self.thresholds.mid_bitrate {
        return Some(StreamQuality::Mid);
    }
    if self.current_stats.packet_loss < self.thresholds.high_loss && 
       self.current_stats.bitrate_bps > self.thresholds.high_bitrate {
        return Some(StreamQuality::High);
    }
    None
}

После обновления статистики сети и определения оптимального уровня качества, реализуем следующую логику:

Инициализация и защита от дребезга (Debounce)
  • Первичный запуск: Если текущее качество еще не задано, система сразу выставляет расчетное значение, отправляет пользователю команду SwitchQualityLayer и завершает шаг.

  • Сравнение состояний: Если текущее качество уже известно, сравниваем его с новым расчетным значением:

    • Качество упало (Ordering::Less): Мы реагируем мгновенно. Снижение качества происходит без задержек, чтобы пользователь не столкнулся с буферизацией. Счетчик стабильных сигналов при этом обнуляется.

    • Качество не изменилось (Ordering::Equal): Переключение не требуется.

    • Качество выросло (Ordering::Greater): Здесь включается защита от ложных срабатываний. Чтобы не повышать качество на основе краткосрочного скачка сети, нам нужно ввести счетчик consecutive_high_signals. Качество вырастет только тогда, когда мы получим подряд больше стабильных сигналов, чем указано в константе REQUIRED_STABLE_SIGNALS. До этого момента счетчик просто инкрементируется.

const REQUIRED_STABLE_SIGNALS: usize = 3;
self.update_stats().await;
self.save_stats().await.ok_or_terminate(ctx);
if let Some(quality) = self.get_quality() {
    let current_quality = match self.current_quality {
        Some(current_quality) => current_quality,
        None => {
            self.current_quality = Some(quality);
            self.user
                .send(UserMessage::SwitchQualityLayer { quality })
                .await
                .ok_or_terminate(ctx);
            return;
        }
    };
    let should_switch = match quality.cmp(&current_quality) {
        std::cmp::Ordering::Equal => {
            false
        },
        std::cmp::Ordering::Greater 
            if self.consecutive_high_signals >= REQUIRED_STABLE_SIGNALS => {
            self.consecutive_high_signals = 0;
            true
        }
        std::cmp::Ordering::Greater => {
            self.consecutive_high_signals += 1;
            false
        },
        std::cmp::Ordering::Less => {
            self.consecutive_high_signals = 0;
            true
        }
    };
    if should_switch {
        self.current_quality = Some(quality);
        self.user.send(UserMessage::SwitchQualityLayer { quality })
            .await
            .ok_or_terminate(ctx);
    }
}
9aaf3fbff5d49b16beb56b581f846bdf.pngSimulcast

Перейдем непосредственно к работе со слоями, но в начале дадим немного теории.

Simulcast (одновременная трансляция) — это технология в видеосвязи и стриминге, когда устройство отправляет на сервер несколько копий одного видео, но разного качества и размера (например, слабое, среднее и высокое). Сервер сам выбирает, какой поток отдать каждому зрителю в зависимости от его интернета.

Актор VideoLayerManager

Опишем актор, который будет непосредственно заниматься переключением слоев по сигналу от мониторинга.

pub struct VideoLayerManager {
    pub pc: Arc<RTCPeerConnection>,
    pub peer_id: Uuid,
    pub quality_layers: HashMap<StreamQuality, QualityLayer>,
    pub active_track: Arc<RTCRtpSender>,
    pub track: Arc<TrackLocalStaticRTP>,
    pub packet_forwarder: WeakAddr<VideoPacketForwarder>,
    pub connection_quality: StreamQuality,
    pub active_quality: Option<StreamQuality>,
    pub is_forwarder_running: bool,
}

Актор выполняет три задачи:

  • Создание нового видео трека и подписка на пришедший слой.

  • Переключение слоев после сигнала о смене качества.

  • Чтение обратного RTCP-канала и обработки потерь данных при стриминге.

impl VideoLayerManager {
    pub async fn new(
        pc: Arc<RTCPeerConnection>, 
        peer_id: Uuid, 
        codec: Codec, 
        current_connection_quality: StreamQuality,
        quality: StreamQuality,
        quality_layer: QualityLayer,
    ) -> Result<Self, Error> {
        let mut quality_layers = HashMap::new();
        let track: Arc<TrackLocalStaticRTP> = Arc::new(TrackLocalStaticRTP::new(
            RTCRtpCodecCapability {
                mime_type: codec.to_string(),
                ..Default::default()
            },
            format!("video_{peer_id}"),
            format!("video_{peer_id}")
        ));
        let active_track = pc.add_transceiver_from_track(
                track.clone() as Arc<_>,
                Some(RTCRtpTransceiverInit { direction: RTCRtpTransceiverDirection::Sendonly, send_encodings: vec![] })
            )
            .await?
            .sender()
            .await;
        quality_layers.insert(quality, quality_layer);
        let this = Self { pc, peer_id, track, quality_layers, active_track, active_quality: None, connection_quality: current_connection_quality, packet_forwarder: WeakAddr::default(), is_forwarder_running: false };
        Ok(this) 
    }
}

Поскольку треки могут прийти в каком угодно порядке, запустим первый пришедший.

if let Some((quality, layer)) = self.quality_layers.iter().next() {
    let wake_notification = layer.wake_notification.clone();
    self.spawn_notify_task(addr.clone(), *quality, wake_notification);
    self.active_quality = Some(*quality);
    self.packet_forwarder
        .try_send(VideoPacketForwarderMessage::Start {
            quality: *quality,
            gateway_router: layer.gateway_router.clone()
        })
        .await;
} 

При добавлении новых слоев проверяем нет ли того, который совпадает по качеству с соединением.

VideoLayerManagerMessage::AddLayer { quality, layer } => {
    self.quality_layers.insert(quality, layer.clone());
    self.spawn_notify_task(ctx.addr.clone(), quality, layer.wake_notification);
    if quality == self.connection_quality {
        let message = VideoPacketForwarderMessage::LayerSwitched {
            gateway_router: layer.gateway_router.clone()
        };
        self.active_quality = Some(quality);
        self.packet_forwarder
            .try_send(message)
            .await
            .ok_or_terminate(ctx);
    }
}

При необходимости смены слоя перешлем сигнал актору VideoPacketForwarder (обновленная версия PacketForwarder из прошлой статьи, к нему мы скоро вернемся).

VideoLayerManagerMessage::SwitchQualityLayer { to  } => {
    let Some(layer) = self.quality_layers.get(&to) else { tracing::warn!("Missing layer"); return ; };
    self.active_quality = Some(to);
    self.connection_quality = to;
    self.packet_forwarder
        .try_send(VideoPacketForwarderMessage::LayerSwitched {
            gateway_router: layer.gateway_router.clone()
        })
        .await
        .ok_or_terminate(ctx);
},
Обработка обратной связи по RTCP (PLI и NACK)

Что такое RTCP и зачем он нужен?

RTCP (RTP Control Protocol) — это вспомогательный протокол для RTP (Real-time Transport Protocol). В то время как RTP доставляет сам медиапоток (аудио/видео), RTCP работает в обратном направлении. Он собирает статистику о качестве связи (потери, задержки, джиттер) и передает управляющие команды от получателя к отправителю, позволяя адаптировать стрим под текущие условия сети.

Механизмы борьбы с потерями: PLI и NACK

В реальном времени (WebRTC/стриминг) нельзя использовать обычный TCP с его гарантированной доставкой - это вызовет огромные задержки. Поэтому поверх UDP применяются специализированные механизмы обратной связи:

1. NACK (Negative Acknowledgment) - Выборочный повтор пакетов

  • Что это такое: Сигнал «я потерял конкретные кусочки пазла».

  • Как работает: Получатель замечает дыру в порядковых номерах RTP-пакетов (например, пришли 101, 102, 104 - пакет 103 потерян). Он отправляет обратно NACK-пакет со списком утерянных ID. Отправитель находит их в своем локальном кэше (буфере ретрансляции) и высылает повторно.

  • Плюс: Экономит трафик, точечно восстанавливая мелкие потери без прерывания видео.

2. PLI (Picture Loss Indication) - Запрос нового ключевого кадра

  • Что это такое: Сигнал «я полностью потерял нить повествования, спасайте».

  • Как работает: Если пакетов потеряно слишком много (или потерян критически важный кусок опорного кадра), декодер на стороне клиента больше не может собирать видео. Картинка рассыпается или застывает. Получатель шлет PLI, требуя от источника немедленно сгенерировать и прислать новый полноценный ключевой кадр (Keyframe / I-Frame).

tokio::spawn(async move {
    loop {
        let mut rtcp_buf = [0u8; 1500];
        let Ok((packets, _)) = active_track.read(&mut rtcp_buf).await else { break; };
        for packet in packets {
            if packet.as_any().downcast_ref::<PictureLossIndication>().is_some() {
                let result = addr.send(VideoSubscriptionMessage::ForcePLI).await;
                if result.is_err() {
                    return;
                }
                continue;
            }
            if let Some(nack) = packet.as_any().downcast_ref::<TransportLayerNack>() {
                let mut missing_seqs = Vec::new();
                for pair in &nack.nacks {
                    let mut packet_list = pair.packet_list();
                    missing_seqs.append(&mut packet_list);
                }
                let message = VideoPacketForwarderMessage::MissedPackets(missing_seqs);
                let result = forwarder.send(message).await;
                if result.is_err() {
                    return;
                };
            }
        }
    }
});
Актор KeyframeInterceptor

Опишем актор для отправления сигнала, о том что необходимо сформировать ключевой кадр.

pub struct KeyframeInterceptor {
    pc: Arc<RTCPeerConnection>,
    ssrc: u32,
    last_keyframe_request: Instant,
}

Ssrc - синхронизационный источник медиапотока (Synchronization Source). Это уникальный идентификатор конкретного видеотрека, для которого данный актор перехватывает и генерирует запросы ключевых кадров.

Реализуем безопасную отправку запроса ключевого кадра (PLI) автору трансляции через PeerConnection. Чтобы частые запросы не перегружали сеть тяжелыми I-кадрами, в коде используется временной ограничитель (Rate Limiter), блокирующий повторную отправку чаще раза в две секунды.

async fn send_pli(&mut self) -> Result<(), crate::Error> {
    let since_last_keyframe = self.last_keyframe_request.elapsed();
    if since_last_keyframe < Duration::from_secs(2) { 
        return Ok(());
    }
    self.pc.write_rtcp(&[Box::new(PictureLossIndication {
        media_ssrc: self.ssrc,
        sender_ssrc: 0
    })]).await?;
    self.last_keyframe_request = Instant::now();
    Ok(()) 
}
b96a9e5d41d530aff9a5112d947772a1.pngАктор VideoPacketForwarder

Опишем обновленный транслятор пакетов.

pub struct VideoPacketForwarder {
    track: Arc<TrackLocalStaticRTP>,
    video_layer_manager: Addr<VideoLayerManager>,
    current_quality: Option<StreamQuality>,
    current_channel: Option<Addr<RtpPacketGatewayRouter<Self, VideoRouterContext>>>,
    pending_channel: Option<Addr<RtpPacketGatewayRouter<Self, VideoRouterContext>>>,
    rtp_packet_cache: RtpPacketCache,
    start_instant: Instant,
    current_generated_ts: u32,
    last_sequence_number: u16,
    is_layer_switching: bool,
}

В первой версии наш VideoPacketForwarder работал как прозрачный прокси: он принимал пакеты от RtpPacketGatewayRouter и сразу писал их в трек подписчика. При добавлении Simulcast нам нужно обеспечить бесшовное переключение слоев, а для этого нам придется манипулировать заголовками пакетов, чтобы браузер не сходил с ума от резкой смены timestamp и sequence_number.

Если оставить оригинальные заголовки, то при смене слоя (например, с High на Low) браузер увидит хаотичный прыжок во времени, дропнет декодер и выдаст секундный фриз, судорожно спамя сервер PLI-запросами.

Вместо того чтобы городить тяжелые буферы на бэкенде и высматривать первый байт ключевого кадра в условиях асинхронной гонки потоков, мы сделаем две вещи:

  1. Мгновенный свитч: переключаем рельсы на новый слой атомарно по первому же прилетевшему пакету. Нам плевать, дельта-кадр это или нет, если декодеру клиента станет плохо, он сам мгновенно пришлет нам PLI через ForcePLI, мы прокинем его издателю, и за один RTT картинка восстановится без замерзания экрана.

  2. Монотонный генератор таймстампов: мы полностью выкинем оригинальное время отправителя и введем единую шкалу времени на основе монотонного таймера бэкенда.

async fn forward(
  &mut self, 
  ctx: &Ctx<'_, Self>, 
  quality: StreamQuality, 
  mut packet: Packet
) -> Result<(), Error> {
    if self.is_layer_switching && self.current_quality != Some(quality) {
        self.current_quality = Some(quality);
        self.is_layer_switching = false;
        if let Some(pending_channel) = self.pending_channel.take() {
            if let Some(cancelled_channel) = self.current_channel.replace(pending_channel) {
                let _ = cancelled_channel.do_send(RtpPacketGatewayRouterMessage::Unsubscribe(ctx.addr.clone()));
            }
        }
    }
    if self.current_quality == Some(quality) {
        self.modify_header(&mut packet);
        self.write_packet(packet).await?;
    }
    Ok(())
}
0533e3bc5214d964523ffd538de31d34.png

Реализуем модификацию заголовков RTP-пакетов для обеспечения непрерывности медиапотока при переключениях. Код последовательно перезаписывает порядковые номера пакетов и синхронизирует временные метки под единую шкалу времени вещания.

fn modify_header(&mut self, packet: &mut Packet) {
  self.last_sequence_number = self.last_sequence_number.wrapping_add(1);
  packet.header.sequence_number = self.last_sequence_number;
  packet.header.timestamp = self.current_generated_ts;
  if packet.header.marker {
      let elapsed_micros = self.start_instant.elapsed().as_micros() as u64;
        // Пересчитываем timestamp для следующего кадра: микросекунды → тики 90 кГц
      let target_timestamp = elapsed_micros.wrapping_mul(9) / 100;
      self.current_generated_ts = target_timestamp as u32;
  }
}

При сообщении о смене слоя, проверим, что слой уже не меняется и подготовим актор к ожиданию пакетов из нового слоя.

VideoPacketForwarderMessage::LayerSwitched {
    gateway_router: forwarder
} 
if !self.is_layer_switching => {
    self.is_layer_switching = true;
    forwarder
      .do_send(RtpPacketGatewayRouterMessage::Subscribe(ctx.addr.clone()))
      .ok_or_terminate(ctx);
    self.pending_channel = Some(forwarder);
},
9974fdcc061e5f5a18cbe224d91a7ec4.pngДополнительная настройка

Мы добавили акторы для работы со слоями, но ничего не будет работать, пока мы не внесем дополнительные настройки на бэкенд и фронт.

У каждого трека приходящего на сервер есть rid. RTP Stream ID (идентификатор RTP-потока). Это специальная метка в заголовке RTP-пакета, которая помогает браузеру и серверу понять, какому именно потоку или слою видео принадлежит данный пакет.

Чтобы сервер смог правильно прочитать рид, добавим расширение для заголовков.

let mut m = MediaEngine::default();
for uri in [
    "urn:ietf:params:rtp-hdrext:sdes:mid",
    "urn:ietf:params:rtp-hdrext:sdes:rtp-stream-id",
    "urn:ietf:params:rtp-hdrext:sdes:repaired-rtp-stream-id",
] {
    m.register_header_extension(
        RTCRtpHeaderExtensionCapability {
            uri: uri.to_owned(),
        },
            RTPCodecType::Video,
        None,
    )?;
}

m.register_default_codecs()?;

Настроим на фронте addTransceiver для отправки трех независимых слоев видео (low, mid, high) с разными битрейтами, разрешениями и кадровой частотой. С помощью массива sendEncodings браузер параллельно кодирует один видеопоток в разных качествах, позволяя SFU-серверу гибко подстраивать поток под пропускную способность каждого зрителя.

const bitrateSettings = {
    low: 400_000,      
    mid: 2_500_000,   
    high: 4_000_000   
};
const encodings = [
    {
        rid: 'low',
        maxBitrate: bitrateSettings.low,
        scaleResolutionDownBy: getRoundedResolution(BASE_HEIGHT, BASE_WIDTH, 4.0),
        maxFramerate: 24,
    },
    {
        rid: 'mid',
        maxBitrate: bitrateSettings.mid,
        scaleResolutionDownBy: getRoundedResolution(BASE_HEIGHT, BASE_WIDTH, 2.0),
        maxFramerate: 30
    },
    {
        rid: 'high',
        maxBitrate: bitrateSettings.high,
        scaleResolutionDownBy: 1,
        maxFramerate: 30
    }
];
const video_track = video_stream.getVideoTracks()[0];
this.video_transceiver = this.publisher_pc.addTransceiver(video_track, {
    direction: 'sendonly',
    sendEncodings: encodings
});

Чтобы протестировать, что все работает, добавим эмуляцию слабого соединения. К сожалению просто использовать throttling profile не получится, по сколько трансляция видео работает по udp. Но ничего страшного, просто сами ограничим битрейт по выбранному селектору.

const applyNetworkProfile = async () => {
  try {
    const video_sender = app_state.peer_connection.video_transceiver.sender;
    if (!video_sender) return;

    const parameters = video_sender.getParameters();
    
    if (!parameters.encodings || parameters.encodings.length === 0) {
      parameters.encodings = [{}];
    }
    const low_layer = parameters.encodings[0];
    const mid_layer = parameters.encodings[1];
    const high_layer = parameters.encodings[2];

    switch (selectedProfile.value) {
      case 'low':
        low_layer.maxBitrate = 150_000; 
        mid_layer.active = false; 
        high_layer.active = false; 
        break;
      case 'mid':
        low_layer.maxBitrate = 400_000;
        mid_layer.active = true; 
        high_layer.active = false; 
        break;
      case 'original':
      default:
        low_layer.maxBitrate = 400_000;
        mid_layer.active = true; 
        high_layer.active = true; 
        break;
    }
    await video_sender.setParameters(parameters);
  } catch (error) {
    console.error('[WebRTC] Ошибка изменения параметров трека:', error);
  }
};

И когда я добавил эту функцию, все сломалось, ведь я совсем не подумал что произойдет, если пакеты внезапно перестанут приходить из выбранного слоя. Добавим таймауты.

tokio::spawn(async move {
    let timeout_period = Duration::from_millis(500);
    let mut is_timeout_send = false;
    loop {
        let read_rtp_fut = track.read_rtp();
        let fut = timeout(timeout_period, read_rtp_fut);
        let message = match fut.await {
            Ok(Ok((packet, _))) => RtpPacketGatewayRouterMessage::RtpPacket(packet),
            Ok(Err(e)) => {
                tracing::error!("{e}");
                let _ = receiver.terminate().await;
                break;
            },
            Err(_) if is_timeout_send => continue,
            Err(_) => {
                is_timeout_send = true;
                RtpPacketGatewayRouterMessage::Timeout
            }
        };
        let Ok(_) = receiver.do_send(message) else { break; };
    }
});

Теперь при получении сообщения о таймауте, запросим переключения на низкий слой.

VideoPacketForwarderMessage::Timeout => {
    let message = VideoLayerManagerMessage::FallbackToLowQuality;
    self.video_layer_manager.send(message)
        .await
        .ok_or_terminate(ctx);
}
96af03945b12c418a882236579cdccba.png

В случае, если поток все же проснется со временем, то мы оповестим VideoLayerManager, о том что можно вернутся к приоретному слою, через tokio::sync::Notify. В роутере вызовем функцию try_awake.

fn try_awake(&mut self) {
    if self.is_sleeping {
        self.wake_notifier.notify_waiters();
    }
}
RtpPacketGatewayRouterMessage::RtpPacket(packet) => {
                for sub in &self.subscriptions {
                    let message = RtpPacketMessage::Packet(self.context.stream_quality(), packet.clone());
                    sub.do_send(message.into()).ok_or_terminate(ctx);
                }
                self.context.try_awake();
            },
RtpPacketGatewayRouterMessage::RtpPacket(packet) => {
    self.context.try_awake();
    for sub in &self.subscriptions {
        let message = RtpPacketMessage::Packet(
            self.context.stream_quality(), 
            packet.clone()
        );
        sub.do_send(message.into()).ok_or_terminate(ctx);
    }
},

А в VideoLayerManager заспавним таску, которая в цикле будет ожидать оповещения, чтобы переключиться на пробудившийся поток.

fn spawn_notify_task(&self, 
  addr: Addr<Self>, 
  quality: StreamQuality, 
  wake_notification: Arc<Notify> 
) {
    tokio::spawn(async move {
        loop {
            wake_notification.notified().await;
            let message = VideoLayerManagerMessage::LayerAwake { quality };
            let send_fut = addr.send(message).await;
            if send_fut.is_err() {
                break;
            }
        }
    });
}
a9d7efb31fb1d6d71cb890d16edea8b8.pngДемонстрация смены слоев

Демонстрация смены слоев

Ice restart

Однако помимо просто плохого соединения, пользователь может столкнуться также с полной потерей сети или же с ее сменой (например с вайфая на мобильную связь). Для этого случая нам нужно организовать алгоритм восстановления соединения.

ICE Restart — это механизм в технологии WebRTC, который позволяет обновить сетевое соединение между пользователями без разрыва текущего звонка.

Зачем это применяется?

  • Смена сети (Роуминг): Переключение телефона с домашнего Wi-Fi на мобильный интернет (4G/5G) во время разговора.

  • Восстановление связи: Автоматический ремонт соединения, если пакеты данных перестали доходить (статус disconnected или failed).

  • Обновление серверов: Быстрый переход на новые TURN/STUN серверы, если старые перестали отвечать.

Я реализовал инициализацию рестарта на фронте.

this.pc.oniceconnectionstatechange = async (e) => {
        switch (this.pc.iceConnectionState) {
            case 'disconnected': {
                setTimeout(async () => {
                    if (this.pc.iceConnectionState !== "connected") {
                        await this.restart_ice()
                    }
                }, 3000);
                break;
            }
            case 'failed': {
                await this.restart_ice();
                break;
            }
        }
    }
}
async restart_ice() {
    if (this.is_restarting) { return; }
    this.is_restarting = true;
    this.pc.restartIce();
    let offer = await this.pc.createOffer({iceRestart: true});
    await this.pc.setLocalDescription(offer);
    this.ws.send(JSON.stringify({
        kind: 'rtc',
        target: this.target,
        type: 'ice_restart',
        sdp: offer.sdp
    }));
    this.is_restarting = false;
}

Как видно из кода, мы получаем евент о смене статуса соединения, после чего запускаем restartIce, создаем оффер и посылаем его на бек по вебсокету.

На беке нам нужно получить сообщение о необходимости рестарта и в послать обратно ответ.

MessageType::IceRestart { sdp } => {
    let offer_desc = RTCSessionDescription::offer(sdp)?;
    self.pc.set_remote_description(offer_desc).await?;
    let answer = self.pc.create_answer(None).await?;
    self.pc.set_local_description(answer.clone()).await?;
    let message = SignalMessage::Rtc {
        target: Target::Publisher,
        message_type: MessageType::Answer {sdp: answer.sdp }
    };
    self.user.send(UserMessage::SignalMessage(message)).await?;
},
0279d012e4ebb61d3e5b66a652fce2c7.pngЗаключение

Мне очень понравилось делать этот проект, возможно в будущем я сделаю его не просто демкой с видео и аудиосвязью, а полноценным порталом с чатом, ролями, файлообменником и всем прочим, для того чтобы он стал настоящей платформой для общения.

Схожие новости

#Наименование новостиТональностьИнформативностьДата публикации
1Async на Rust завис: как понять, где именно, когда паники нет и стек молчит07.6612-08-2026
2Разработка драйвера сетевого адаптера для Linux. Часть 20830-06-2026
3Как мы приручили vGPU до режима авто без проблем06.322-06-2026
4Свой видеохостинг с P2P и Fediverse: PeerTube08.7612-08-2026
5Куда пропали 14 млрд рынка observability?08.2206-08-2026
6Спутниковая связь в симуляторе NS-3. Часть 208.7103-06-2026
7Управление проектами: 10 самых интересных публикаций за 2 недели010.8307-08-2026
8Микрофронтенды. Стабильная интеграция нескольких SPA-приложений. Часть 20530-06-2026
9Как отговорить себя от написания своей ФС011.6409-08-2026
10[Перевод] Мультиагентные системы как распределенное программное обеспечение08.8629-06-2026

Классификация: Мнения. Схожих патентов: 0. Схожих новостей: 10. Тональность: 0. Информативность: 9.77. Источник: habr.com.