use super::*;
#[test]
fn in_process_bypass_delivers_seq_header_then_gated_keyframe_then_interleaved() {
let mut scheduler = RtmpScheduler::new(10);
let publisher_conn = 100;
let watcher_conn = 2;
assert!(scheduler.new_channel("live".to_string(), publisher_conn));
play(&mut scheduler, watcher_conn, "live");
let watcher_client_id = *scheduler
.connection_to_client_map
.get(&watcher_conn)
.unwrap();
let results = feed_media(&mut scheduler, publisher_conn, 0x09, 0, VIDEO_SEQ);
assert_eq!(
results.len(),
1,
"video sequence header must reach the watcher"
);
assert!(matches!(
&results[0],
ServerResult::OutboundPacket {
is_sequence_header: true,
is_video: true,
..
}
));
let results = feed_media(&mut scheduler, publisher_conn, 0x08, 0, AUDIO_SEQ);
assert_eq!(
results.len(),
1,
"audio sequence header must reach the watcher"
);
assert!(matches!(
&results[0],
ServerResult::OutboundPacket {
is_sequence_header: true,
is_video: false,
..
}
));
assert!(
feed_media(&mut scheduler, publisher_conn, 0x09, 33, DELTA).is_empty(),
"delta before the keyframe must be withheld"
);
assert!(
feed_media(&mut scheduler, publisher_conn, 0x08, 33, AUDIO_FRAME).is_empty(),
"audio before the first keyframe must be withheld"
);
assert!(
!scheduler
.clients
.get(watcher_client_id)
.unwrap()
.has_received_video_keyframe
);
let results = feed_media(&mut scheduler, publisher_conn, 0x09, 66, KEYFRAME);
assert_eq!(results.len(), 1, "the keyframe must be delivered");
assert!(matches!(
&results[0],
ServerResult::OutboundPacket {
is_keyframe: true,
is_video: true,
..
}
));
assert!(
scheduler
.clients
.get(watcher_client_id)
.unwrap()
.has_received_video_keyframe
);
let results = feed_media(&mut scheduler, publisher_conn, 0x08, 70, AUDIO_FRAME);
assert_eq!(results.len(), 1, "audio flows after the keyframe");
assert!(matches!(
&results[0],
ServerResult::OutboundPacket {
is_video: false,
is_keyframe: false,
..
}
));
let results = feed_media(&mut scheduler, publisher_conn, 0x09, 99, DELTA);
assert_eq!(results.len(), 1, "delta flows after the keyframe");
assert!(matches!(
&results[0],
ServerResult::OutboundPacket {
is_video: true,
is_keyframe: false,
..
}
));
}
#[test]
fn in_process_bypass_populates_seq_headers_and_gop_cache_for_replay() {
let mut scheduler = RtmpScheduler::new(10);
let publisher_conn = 100;
assert!(scheduler.new_channel("live".to_string(), publisher_conn));
feed_media(&mut scheduler, publisher_conn, 0x09, 0, VIDEO_SEQ);
feed_media(&mut scheduler, publisher_conn, 0x08, 0, AUDIO_SEQ);
feed_media(&mut scheduler, publisher_conn, 0x09, 33, KEYFRAME);
feed_media(&mut scheduler, publisher_conn, 0x08, 34, AUDIO_FRAME);
feed_media(&mut scheduler, publisher_conn, 0x09, 66, DELTA);
feed_media(&mut scheduler, publisher_conn, 0x09, 99, KEYFRAME);
let channel = scheduler.channels.get("live").unwrap();
assert_eq!(channel.video_sequence_header.as_deref(), Some(VIDEO_SEQ));
assert_eq!(channel.video_timestamp, RtmpTimestamp { value: 0 });
assert_eq!(channel.audio_sequence_header.as_deref(), Some(AUDIO_SEQ));
assert_eq!(channel.audio_timestamp, RtmpTimestamp { value: 0 });
assert!(
channel.gops.frozen_count() >= 1,
"the completed GOP must be frozen for replay"
);
}
#[test]
fn fanout_skips_watchers_whose_action_points_at_another_stream() {
let mut scheduler = RtmpScheduler::new(10);
let publisher_conn = 100;
let watcher_conn = 2;
assert!(scheduler.new_channel("stream_a".to_string(), publisher_conn));
play(&mut scheduler, watcher_conn, "stream_a");
let client_id = *scheduler
.connection_to_client_map
.get(&watcher_conn)
.unwrap();
scheduler.clients.get_mut(client_id).unwrap().current_action = ClientAction::Watching {
stream_key: "stream_b".to_string(),
stream_id: 1,
};
let results = feed_video(&mut scheduler, "stream_a", 0, KEYFRAME);
assert!(
results.is_empty(),
"frames must not be delivered through a stale watcher membership"
);
assert!(
!scheduler
.clients
.get(client_id)
.unwrap()
.has_received_video_keyframe,
"a mismatched channel must not open the client's keyframe gate"
);
}