openrtc 1.0.1

OpenRTC: a Rust-first P2P runtime for device discovery, signaling, and iroh/QUIC networking.
Documentation
/// Wire format for heartbeat ping/pong frames.
///
/// Framed iroh ping payload (17 bytes):
///   [0..8]  u64 BE  sequence number (monotonic per sender per transport)
///   [8..16] i64 BE  sender_unix_ms
///   [16]    u8      flags (bit 0: must-pong)
///
/// Framed iroh pong payload (25 bytes) = ping bytes + 8 more:
///   [17..25] i64 BE responder_recv_unix_ms
///
/// WebRTC data-channel payloads intentionally omit the flags byte so they
/// match the browser TypeScript heartbeat codec:
///   ping: [seq u64 BE][sender_unix_ms i64 BE]
///   pong: [seq u64 BE][sender_unix_ms i64 BE][recv_unix_ms i64 BE]
///
/// On framed iroh path: wrapped in [8-byte magic][u32 len][u8 typeId][payload]
/// On WebRTC data channel: prefix byte 0x03 (ping) or 0x04 (pong) + payload

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; // must-pong
    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> {
    // Accept both iroh-framed (17-byte payload with flags) and WebRTC
    // (16-byte payload without flags) heartbeat pings.
    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()?),
    })
}

/// Build a complete framed ping for iroh streams: [u32 len BE][typeId=3][payload]
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)
}

/// Build a complete framed pong for iroh streams: [u32 len BE][typeId=4][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
}

/// Build a WebRTC data-channel ping: [0x03][payload]
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
}

/// Build a WebRTC data-channel pong: [0x04][payload]
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);
    }
}