use super::{
new_media_serializer, serialize_media, Client, ClientAction, MediaChannel,
ReceivedDataType, RtmpScheduler, ServerResult, OUTBOUND_CHUNK_SIZE,
};
use crate::flv::flv_tag_body::{
is_audio_sequence_header, is_video_keyframe, is_video_sequence_header,
};
use crate::rtmp::gop::FrameData;
use crate::rtmp::write_queue::QUEUE_WARN_BYTES;
use log::debug;
use rml_rtmp::sessions::ServerSessionResult;
const JOIN_REPLAY_HEADROOM: usize = 64 * 1024;
pub(super) const JOIN_REPLAY_BUDGET_BYTES: usize = QUEUE_WARN_BYTES - JOIN_REPLAY_HEADROOM;
pub(super) const MSG_HEADER_MAX: usize = 18;
pub(super) const CONT_HEADER_MAX: usize = 3;
pub(super) fn gop_wire_size(payload: usize, frame_count: usize) -> usize {
let per_frame = frame_count.saturating_mul(MSG_HEADER_MAX);
let continuation = payload
.div_ceil(OUTBOUND_CHUNK_SIZE)
.saturating_mul(CONT_HEADER_MAX);
payload
.saturating_add(per_frame)
.saturating_add(continuation)
}
impl RtmpScheduler {
pub(super) fn advance_serving_prefix(&mut self, server_results: &[ServerResult], target: usize) -> usize {
while self.serving_prefix_scan_pos < server_results.len() {
if let ServerResult::OutboundPacket {
target_connection_id,
bytes,
..
} = &server_results[self.serving_prefix_scan_pos]
{
if *target_connection_id == target {
self.serving_prefix_bytes =
self.serving_prefix_bytes.saturating_add(bytes.len());
}
}
self.serving_prefix_scan_pos += 1;
}
self.serving_prefix_bytes
}
pub(super) fn handle_play_requested(
&mut self,
requested_connection_id: usize,
request_id: u32,
app_name: String,
stream_key: String,
stream_id: u32,
server_results: &mut Vec<ServerResult>,
) {
debug!("Rtmp play requested on app '{app_name}' and stream key '{stream_key}'");
let previous_watch = {
let client_id = self
.connection_to_client_map
.get(&requested_connection_id)
.unwrap();
let client = self.clients.get(*client_id).unwrap();
match &client.current_action {
ClientAction::Watching {
stream_key: old_stream_key,
..
} if *old_stream_key != stream_key => Some((*client_id, old_stream_key.clone())),
_ => None,
}
};
if let Some((client_id, old_stream_key)) = previous_watch {
self.play_ended(client_id, old_stream_key);
}
let serving_backlog = self.serving_connection_backlog_bytes;
let same_batch_prefix =
self.advance_serving_prefix(server_results.as_slice(), requested_connection_id);
let accept_result;
let mut join_burst = Vec::new();
{
let client_id = self
.connection_to_client_map
.get(&requested_connection_id)
.unwrap();
let client = self.clients.get_mut(*client_id).unwrap();
client.current_action = ClientAction::Watching {
stream_key: stream_key.clone(),
stream_id,
};
client.has_received_video_keyframe = false;
let channel = self
.channels
.entry(stream_key.clone())
.or_insert_with(|| MediaChannel::new(self.gop_limit));
channel.watching_client_ids.insert(*client_id);
accept_result = client.session.accept_request(request_id);
if let Ok(ref accept_results) = accept_result {
let accept_prefix_bytes =
join_replay_prefix_bytes(serving_backlog, same_batch_prefix, accept_results);
build_join_burst(
channel,
client,
requested_connection_id,
stream_id,
accept_prefix_bytes,
&mut join_burst,
);
}
}
match accept_result {
Err(error) => {
debug!(
"Rtmp client error occurred accepting playback request: {:?}",
error
);
server_results.push(ServerResult::DisconnectConnection {
connection_id: requested_connection_id,
});
return;
}
Ok(results) => {
self.handle_session_results(requested_connection_id, results, server_results);
server_results.extend(join_burst);
}
}
}
}
pub(super) fn select_replay_start(sizes: &[usize], budget: usize) -> usize {
let mut total: usize = 0;
let mut start = sizes.len();
for (i, &size) in sizes.iter().enumerate().rev() {
match total.checked_add(size) {
Some(sum) if sum <= budget => {
total = sum;
start = i;
}
_ => break,
}
}
start
}
pub(super) fn join_replay_prefix_bytes(
connection_backlog: usize,
same_batch_prefix: usize,
accept_results: &[ServerSessionResult],
) -> usize {
let accept_packet_bytes = accept_results
.iter()
.fold(0usize, |acc, result| match result {
ServerSessionResult::OutboundResponse(packet) => acc.saturating_add(packet.bytes.len()),
_ => acc,
});
connection_backlog
.saturating_add(same_batch_prefix)
.saturating_add(accept_packet_bytes)
}
pub(super) fn build_join_burst(
channel: &MediaChannel,
client: &mut Client,
connection_id: usize,
stream_id: u32,
accept_prefix_bytes: usize,
out: &mut Vec<ServerResult>,
) {
let mut burst_serializer = new_media_serializer();
let mut prefix_wire_bytes = accept_prefix_bytes;
if let Some(ref metadata) = channel.metadata {
match client.session.send_metadata(stream_id, metadata) {
Ok(packet) => {
prefix_wire_bytes = prefix_wire_bytes.saturating_add(packet.bytes.len());
out.push(ServerResult::outbound(
connection_id,
packet,
false,
false,
false,
));
}
Err(error) => {
debug!(
"Rtmp client error occurred sending existing metadata to new client: {:?}",
error
);
out.push(ServerResult::DisconnectConnection { connection_id });
return;
}
}
}
if let Some(ref data) = channel.video_sequence_header {
match serialize_media(
&mut burst_serializer,
ReceivedDataType::Video,
stream_id,
data.clone(),
channel.video_timestamp,
false,
) {
Ok(packet) => {
prefix_wire_bytes = prefix_wire_bytes.saturating_add(packet.bytes.len());
out.push(ServerResult::outbound(
connection_id,
packet,
false,
true,
true,
));
}
Err(error) => {
debug!(
"Rtmp client error occurred sending video header to new client: {:?}",
error
);
out.push(ServerResult::DisconnectConnection { connection_id });
return;
}
}
}
if let Some(ref data) = channel.audio_sequence_header {
match serialize_media(
&mut burst_serializer,
ReceivedDataType::Audio,
stream_id,
data.clone(),
channel.audio_timestamp,
false,
) {
Ok(packet) => {
prefix_wire_bytes = prefix_wire_bytes.saturating_add(packet.bytes.len());
out.push(ServerResult::outbound(
connection_id,
packet,
false,
true,
false,
));
}
Err(error) => {
debug!(
"Rtmp client error occurred sending audio header to new client: {:?}",
error
);
out.push(ServerResult::DisconnectConnection { connection_id });
return;
}
}
}
let budget = JOIN_REPLAY_BUDGET_BYTES.saturating_sub(prefix_wire_bytes);
let frozen: Vec<_> = channel.gops.get_frozen_gops().collect();
let current_frames = channel.gops.current_frames();
let mut sizes: Vec<usize> = frozen
.iter()
.map(|gop| gop_wire_size(gop.byte_size(), gop.frame_count()))
.collect();
sizes.push(gop_wire_size(
channel.gops.current_byte_size(),
current_frames.len(),
));
let start = select_replay_start(&sizes, budget);
let mut replayed_keyframe = false;
for segment_index in start..sizes.len() {
let frames = if segment_index < frozen.len() {
frozen[segment_index].frames()
} else {
current_frames
};
if !replayed_keyframe && gop_contains_video_keyframe(frames) {
replayed_keyframe = true;
client.has_received_video_keyframe = true;
}
for frame_data in frames {
match frame_data {
FrameData::Video { timestamp, data } => {
if !replayed_keyframe {
continue;
}
let is_keyframe = is_video_keyframe(data);
let is_sequence_header = is_video_sequence_header(data);
match serialize_media(
&mut burst_serializer,
ReceivedDataType::Video,
stream_id,
data.clone(),
*timestamp,
false,
) {
Ok(packet) => out.push(ServerResult::outbound(
connection_id,
packet,
is_keyframe,
is_sequence_header,
true,
)),
Err(error) => {
debug!(
"Rtmp client error occurred sending video data to new client: {:?}",
error
);
out.push(ServerResult::DisconnectConnection { connection_id });
return;
}
}
}
FrameData::Audio { timestamp, data } => {
let is_sequence_header = is_audio_sequence_header(data);
match serialize_media(
&mut burst_serializer,
ReceivedDataType::Audio,
stream_id,
data.clone(),
*timestamp,
false,
) {
Ok(packet) => out.push(ServerResult::outbound(
connection_id,
packet,
false,
is_sequence_header,
false,
)),
Err(error) => {
debug!(
"Rtmp client error occurred sending audio data to new client: {:?}",
error
);
out.push(ServerResult::DisconnectConnection { connection_id });
return;
}
}
}
}
}
}
}
pub(super) fn gop_contains_video_keyframe(frames: &[FrameData]) -> bool {
frames.iter().any(|f| match f {
FrameData::Video { data, .. } => is_video_keyframe(data),
FrameData::Audio { .. } => false,
})
}