Skip to content

Commit c8a8cf6

Browse files
committed
feat: merge back an unfortunate amount of features
This includes: - MediaStream stream & track id (msid attribute) - Retransmissions using rtx payload type - NACKs to support retransmissions - A basic transport-cc framework to build upon
1 parent 3c01f7f commit c8a8cf6

32 files changed

Lines changed: 1916 additions & 375 deletions

Cargo.toml

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,3 +32,6 @@ log = "0.4"
3232

3333
[workspace.lints.rust]
3434
unreachable_pub = "warn"
35+
36+
[patch.crates-io]
37+
rtcp-types = { git = "https://github.com/kbalt/rtcp-types.git", rev = "dffb2bead61d5bf387fbc031f1c2c9036455f7c3" }

examples/make_call.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -58,7 +58,7 @@ async fn main() -> Result<(), Box<dyn std::error::Error>> {
5858
.unwrap();
5959

6060
// Add an audio stream using the previously defined codecs
61-
sdp_session.add_media(audio, Direction::SendRecv);
61+
sdp_session.add_media(audio, Direction::SendRecv, None, None);
6262

6363
let mut outbound_call = registration
6464
.make_call(

media/rtc/Cargo.toml

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,7 @@ slotmap = "1.0.7"
2727
thiserror = "2"
2828
time = "0.3"
2929

30-
tokio = { version = "1", features = [
31-
"net",
32-
"macros",
33-
"time",
34-
], default-features = false, optional = true }
30+
tokio = { version = "1", features = ["net", "macros", "time"], optional = true }
3531
quinn-udp = { version = "0.5", optional = true }
3632
local-ip-address = { version = "0.6", optional = true }
3733

@@ -40,3 +36,7 @@ default = ["tokio"]
4036
tokio = ["dep:tokio", "dep:quinn-udp", "dep:local-ip-address"]
4137

4238
vendor-openssl = ["openssl/vendored"]
39+
40+
[dev-dependencies]
41+
tokio = { version = "1", features = ["rt"] }
42+
env_logger = "0.11"

media/rtc/examples/echo.rs

Lines changed: 127 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,127 @@
1+
use std::{
2+
net::Ipv4Addr,
3+
time::{Duration, Instant},
4+
};
5+
6+
use ezk_rtc::{
7+
Mtu, OpenSslContext,
8+
rtp_session::SendRtpPacket,
9+
sdp::{
10+
BundlePolicy, Codec, Codecs, LocalMediaId, RtcpMuxPolicy, SdpSession, SdpSessionConfig,
11+
SdpSessionEvent, TransportType,
12+
},
13+
tokio::TokioIoState,
14+
};
15+
use sdp_types::{Direction, MediaType, SessionDescription};
16+
17+
pub(crate) fn make_session(config: SdpSessionConfig) -> (LocalMediaId, SdpSession) {
18+
let mut session = SdpSession::new(
19+
OpenSslContext::try_new().unwrap(),
20+
Ipv4Addr::LOCALHOST.into(),
21+
config,
22+
);
23+
24+
let audio = session
25+
.add_local_media(
26+
Codecs::new(MediaType::Video).with_codec(Codec::VP8.with_rtx()),
27+
Direction::SendRecv,
28+
)
29+
.unwrap();
30+
31+
(audio, session)
32+
}
33+
34+
#[tokio::main(flavor = "current_thread")]
35+
async fn main() {
36+
env_logger::builder().is_test(true).init();
37+
38+
let (local_media_id, mut sdp_session) = make_session(SdpSessionConfig {
39+
offer_transport: TransportType::DtlsSrtp,
40+
offer_ice: true,
41+
offer_avpf: true,
42+
rtcp_mux_policy: RtcpMuxPolicy::Require,
43+
bundle_policy: BundlePolicy::MaxBundle,
44+
mtu: Mtu::new(1400),
45+
});
46+
47+
let mut io = TokioIoState::new_with_local_ips().unwrap();
48+
49+
sdp_session.add_media(local_media_id, Direction::SendOnly, None, None);
50+
51+
io.handle_transport_changes(&mut sdp_session).await.unwrap();
52+
53+
println!("Paste SDP offer:");
54+
let mut offer = String::new();
55+
while !offer.ends_with("\n\n") {
56+
std::io::stdin().read_line(&mut offer).unwrap();
57+
}
58+
let offer = SessionDescription::parse(&offer.into()).unwrap();
59+
let answer = sdp_session.receive_sdp_offer(offer).unwrap();
60+
io.handle_transport_changes(&mut sdp_session).await.unwrap();
61+
let answer = sdp_session.create_sdp_answer(answer);
62+
63+
println!("SDP Answer:\n{answer}");
64+
65+
let base_time = Instant::now();
66+
67+
loop {
68+
while let Ok(event) = io.poll_session(&mut sdp_session).await {
69+
handle_event(base_time, &mut io, &mut sdp_session, event);
70+
}
71+
}
72+
}
73+
74+
fn handle_event(
75+
base_time: Instant,
76+
io: &mut TokioIoState,
77+
sdp_session: &mut SdpSession,
78+
event: SdpSessionEvent,
79+
) {
80+
match event {
81+
SdpSessionEvent::MediaAdded(e) => {
82+
println!("{e:?}");
83+
}
84+
SdpSessionEvent::MediaChanged(e) => {
85+
println!("{e:?}");
86+
}
87+
SdpSessionEvent::MediaRemoved(e) => {
88+
println!("{e:?}");
89+
}
90+
SdpSessionEvent::IceGatheringState(e) => {
91+
println!("{e:?}");
92+
}
93+
SdpSessionEvent::IceConnectionState(e) => {
94+
println!("{e:?}");
95+
}
96+
SdpSessionEvent::TransportConnectionState(e) => {
97+
println!("{e:?}");
98+
}
99+
SdpSessionEvent::SendData {
100+
transport_id,
101+
component,
102+
data,
103+
source,
104+
target,
105+
} => {
106+
io.send(transport_id, component, data, source, target);
107+
}
108+
SdpSessionEvent::ReceiveRTP {
109+
media_id,
110+
rtp_packet,
111+
} => {
112+
let timestamp = Duration::from_secs_f64(rtp_packet.timestamp.0 as f64 / 90_000.0);
113+
114+
let rtp_time = base_time + timestamp;
115+
116+
let mut outbound_media = sdp_session.outbound_media(media_id).unwrap();
117+
118+
outbound_media.send_rtp(
119+
SendRtpPacket::new(rtp_time, rtp_packet.pt, rtp_packet.payload)
120+
.marker(rtp_packet.marker)
121+
.send_at(base_time),
122+
);
123+
}
124+
SdpSessionEvent::ReceivePictureLossIndication { .. } => {}
125+
SdpSessionEvent::ReceiveFullIntraRefresh { .. } => {}
126+
}
127+
}

media/rtc/examples/roundtrip.rs

Lines changed: 144 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,144 @@
1+
use std::{
2+
net::Ipv4Addr,
3+
time::{Duration, Instant},
4+
};
5+
6+
use bytes::Bytes;
7+
use ezk_rtc::{
8+
OpenSslContext,
9+
rtp_session::SendRtpPacket,
10+
sdp::{
11+
Codec, Codecs, LocalMediaId, SdpSession, SdpSessionConfig, SdpSessionEvent, TransportType,
12+
},
13+
tokio::TokioIoState,
14+
};
15+
use sdp_types::{Direction, MediaType};
16+
use tokio::{select, time::interval};
17+
18+
pub(crate) fn make_session(config: SdpSessionConfig) -> (LocalMediaId, SdpSession) {
19+
let mut session = SdpSession::new(
20+
OpenSslContext::try_new().unwrap(),
21+
Ipv4Addr::LOCALHOST.into(),
22+
config,
23+
);
24+
25+
let audio = session
26+
.add_local_media(
27+
Codecs::new(MediaType::Video).with_codec(Codec::H264.with_rtx()),
28+
Direction::SendRecv,
29+
)
30+
.unwrap();
31+
32+
(audio, session)
33+
}
34+
35+
#[tokio::main(flavor = "current_thread")]
36+
async fn main() {
37+
env_logger::builder().is_test(true).init();
38+
39+
let (local_media_id1, mut sdp_session1) = make_session(SdpSessionConfig {
40+
offer_ice: true,
41+
offer_transport: TransportType::Rtp,
42+
offer_avpf: true,
43+
..Default::default()
44+
});
45+
let (_local_media_id2, mut sdp_session2) = make_session(SdpSessionConfig {
46+
offer_transport: TransportType::Rtp,
47+
..Default::default()
48+
});
49+
50+
let mut io1 = TokioIoState::new_with_local_ips().unwrap();
51+
let mut io2 = TokioIoState::new_with_local_ips().unwrap();
52+
53+
sdp_session1.add_media(local_media_id1, Direction::SendOnly, None, None);
54+
55+
io1.handle_transport_changes(&mut sdp_session1)
56+
.await
57+
.unwrap();
58+
59+
let offer = sdp_session1.create_sdp_offer();
60+
61+
println!("Offer:\n{offer}");
62+
63+
let answer = sdp_session2.receive_sdp_offer(offer).unwrap();
64+
io2.handle_transport_changes(&mut sdp_session2)
65+
.await
66+
.unwrap();
67+
let answer = sdp_session2.create_sdp_answer(answer);
68+
69+
println!("Answer:\n{answer}");
70+
71+
sdp_session1.receive_sdp_answer(answer).unwrap();
72+
io1.handle_transport_changes(&mut sdp_session1)
73+
.await
74+
.unwrap();
75+
76+
let mut send_interval = interval(Duration::from_millis(16));
77+
78+
let payload = Bytes::from(vec![0u8; 1300]);
79+
80+
loop {
81+
select! {
82+
event = io1.poll_session(&mut sdp_session1) => {
83+
handle_event(&mut io1, event.unwrap());
84+
continue;
85+
}
86+
event = io2.poll_session(&mut sdp_session2) => {
87+
handle_event(&mut io2, event.unwrap());
88+
continue;
89+
}
90+
_ = send_interval.tick() => {
91+
// fallthrough
92+
}
93+
}
94+
95+
let media = sdp_session1.media_iter().next().unwrap();
96+
let media_id = media.id();
97+
98+
let now = Instant::now();
99+
100+
if let Some(mut outbound) = sdp_session1.outbound_media(media_id) {
101+
for _ in 0..1 {
102+
outbound.send_rtp(SendRtpPacket::new(now, 96, payload.clone()));
103+
}
104+
}
105+
}
106+
}
107+
108+
fn handle_event(io: &mut TokioIoState, event: SdpSessionEvent) {
109+
match event {
110+
SdpSessionEvent::MediaAdded(e) => {
111+
println!("{e:?}");
112+
}
113+
SdpSessionEvent::MediaChanged(e) => {
114+
println!("{e:?}");
115+
}
116+
SdpSessionEvent::MediaRemoved(e) => {
117+
println!("{e:?}");
118+
}
119+
SdpSessionEvent::IceGatheringState(e) => {
120+
println!("{e:?}");
121+
}
122+
SdpSessionEvent::IceConnectionState(e) => {
123+
println!("{e:?}");
124+
}
125+
SdpSessionEvent::TransportConnectionState(e) => {
126+
println!("{e:?}");
127+
}
128+
SdpSessionEvent::SendData {
129+
transport_id,
130+
component,
131+
data,
132+
source,
133+
target,
134+
} => {
135+
io.send(transport_id, component, data, source, target);
136+
}
137+
SdpSessionEvent::ReceiveRTP {
138+
media_id: _,
139+
rtp_packet: _,
140+
} => {}
141+
SdpSessionEvent::ReceivePictureLossIndication { .. } => {}
142+
SdpSessionEvent::ReceiveFullIntraRefresh { .. } => {}
143+
}
144+
}

0 commit comments

Comments
 (0)