pub const IROH_CONTROL_MAGIC: &[u8; 8] = b"ORTCCTL1";
pub const IROH_CONTROL_HEADER_LEN: usize = IROH_CONTROL_MAGIC.len() + 5;
pub const TYPE_PING: u8 = 3;
pub const TYPE_PONG: u8 = 4;
pub const TYPE_MANUAL_DISCONNECT: u8 = 5;
pub const FRAMED_PING_LEN: usize = 17;
pub const FRAMED_PONG_LEN: usize = 25;
pub const WEBRTC_PING_LEN: usize = 16;
pub const WEBRTC_PONG_LEN: usize = 24;
#[derive(Debug, Clone)]
pub struct Ping {
pub seq: u64,
pub sender_unix_ms: i64,
}
#[derive(Debug, Clone)]
pub struct Pong {
pub seq: u64,
pub sender_unix_ms: i64,
pub recv_unix_ms: i64,
}
pub fn encode_ping(seq: u64, sender_unix_ms: i64) -> [u8; FRAMED_PING_LEN] {
let mut buf = [0u8; FRAMED_PING_LEN];
buf[0..8].copy_from_slice(&seq.to_be_bytes());
buf[8..16].copy_from_slice(&sender_unix_ms.to_be_bytes());
buf[16] = 0x01; buf
}
pub fn encode_pong(ping: &Ping, recv_unix_ms: i64) -> [u8; FRAMED_PONG_LEN] {
let mut buf = [0u8; FRAMED_PONG_LEN];
buf[0..8].copy_from_slice(&ping.seq.to_be_bytes());
buf[8..16].copy_from_slice(&ping.sender_unix_ms.to_be_bytes());
buf[16] = 0x00;
buf[17..25].copy_from_slice(&recv_unix_ms.to_be_bytes());
buf
}
pub fn decode_ping(payload: &[u8]) -> Option<Ping> {
if payload.len() < WEBRTC_PING_LEN {
return None;
}
Some(Ping {
seq: u64::from_be_bytes(payload[0..8].try_into().ok()?),
sender_unix_ms: i64::from_be_bytes(payload[8..16].try_into().ok()?),
})
}
pub fn decode_pong(payload: &[u8]) -> Option<Pong> {
if payload.len() < WEBRTC_PONG_LEN {
return None;
}
let recv_offset = if payload.len() >= FRAMED_PONG_LEN {
17
} else {
16
};
Some(Pong {
seq: u64::from_be_bytes(payload[0..8].try_into().ok()?),
sender_unix_ms: i64::from_be_bytes(payload[8..16].try_into().ok()?),
recv_unix_ms: i64::from_be_bytes(payload[recv_offset..recv_offset + 8].try_into().ok()?),
})
}
pub fn build_framed_ping(seq: u64, sender_unix_ms: i64) -> Vec<u8> {
let payload = encode_ping(seq, sender_unix_ms);
build_frame(TYPE_PING, &payload)
}
pub fn build_framed_pong(ping: &Ping, recv_unix_ms: i64) -> Vec<u8> {
let payload = encode_pong(ping, recv_unix_ms);
build_frame(TYPE_PONG, &payload)
}
pub fn build_framed_manual_disconnect() -> Vec<u8> {
build_frame(TYPE_MANUAL_DISCONNECT, &[])
}
fn build_frame(type_id: u8, payload: &[u8]) -> Vec<u8> {
let mut frame = Vec::with_capacity(IROH_CONTROL_HEADER_LEN + payload.len());
frame.extend_from_slice(IROH_CONTROL_MAGIC);
let len = (1 + payload.len()) as u32;
frame.extend_from_slice(&len.to_be_bytes());
frame.push(type_id);
frame.extend_from_slice(payload);
frame
}
pub fn build_webrtc_ping(seq: u64, sender_unix_ms: i64) -> Vec<u8> {
let mut frame = Vec::with_capacity(1 + WEBRTC_PING_LEN);
frame.push(TYPE_PING);
frame.extend_from_slice(&seq.to_be_bytes());
frame.extend_from_slice(&sender_unix_ms.to_be_bytes());
frame
}
pub fn build_webrtc_pong(ping: &Ping, recv_unix_ms: i64) -> Vec<u8> {
let mut frame = Vec::with_capacity(1 + WEBRTC_PONG_LEN);
frame.push(TYPE_PONG);
frame.extend_from_slice(&ping.seq.to_be_bytes());
frame.extend_from_slice(&ping.sender_unix_ms.to_be_bytes());
frame.extend_from_slice(&recv_unix_ms.to_be_bytes());
frame
}
pub fn now_unix_ms() -> i64 {
#[cfg(target_arch = "wasm32")]
{
js_sys::Date::now() as i64
}
#[cfg(not(target_arch = "wasm32"))]
{
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ping_roundtrip() {
let frame = build_framed_ping(42, 1_700_000_000_000);
assert_eq!(&frame[..IROH_CONTROL_MAGIC.len()], IROH_CONTROL_MAGIC);
assert_eq!(frame[IROH_CONTROL_MAGIC.len() + 4], TYPE_PING);
let payload = &frame[IROH_CONTROL_HEADER_LEN..];
let ping = decode_ping(payload).unwrap();
assert_eq!(ping.seq, 42);
assert_eq!(ping.sender_unix_ms, 1_700_000_000_000);
}
#[test]
fn pong_roundtrip() {
let ping = Ping {
seq: 7,
sender_unix_ms: 100,
};
let frame = build_framed_pong(&ping, 200);
assert_eq!(&frame[..IROH_CONTROL_MAGIC.len()], IROH_CONTROL_MAGIC);
assert_eq!(frame[IROH_CONTROL_MAGIC.len() + 4], TYPE_PONG);
let payload = &frame[IROH_CONTROL_HEADER_LEN..];
let pong = decode_pong(payload).unwrap();
assert_eq!(pong.seq, 7);
assert_eq!(pong.sender_unix_ms, 100);
assert_eq!(pong.recv_unix_ms, 200);
}
#[test]
fn webrtc_ping_prefix() {
let frame = build_webrtc_ping(1, 0);
assert_eq!(frame[0], TYPE_PING);
assert_eq!(frame.len(), 1 + WEBRTC_PING_LEN);
let ping = decode_ping(&frame[1..]).unwrap();
assert_eq!(ping.seq, 1);
}
#[test]
fn webrtc_pong_prefix() {
let ping = Ping {
seq: 3,
sender_unix_ms: 50,
};
let frame = build_webrtc_pong(&ping, 75);
assert_eq!(frame[0], TYPE_PONG);
assert_eq!(frame.len(), 1 + WEBRTC_PONG_LEN);
let pong = decode_pong(&frame[1..]).unwrap();
assert_eq!(pong.seq, 3);
assert_eq!(pong.recv_unix_ms, 75);
}
#[test]
fn decode_ping_accepts_webrtc_payload_without_flags() {
let mut payload = Vec::new();
payload.extend_from_slice(&9u64.to_be_bytes());
payload.extend_from_slice(&123i64.to_be_bytes());
let ping = decode_ping(&payload).unwrap();
assert_eq!(ping.seq, 9);
assert_eq!(ping.sender_unix_ms, 123);
}
#[test]
fn decode_pong_accepts_webrtc_payload_without_flags() {
let mut payload = Vec::new();
payload.extend_from_slice(&11u64.to_be_bytes());
payload.extend_from_slice(&222i64.to_be_bytes());
payload.extend_from_slice(&333i64.to_be_bytes());
let pong = decode_pong(&payload).unwrap();
assert_eq!(pong.seq, 11);
assert_eq!(pong.sender_unix_ms, 222);
assert_eq!(pong.recv_unix_ms, 333);
}
}