use base64::Engine as _;
use serde::{Deserialize, Serialize};
fn parse_application_key_agreement_public_key(
json: &serde_json::Value,
) -> Option<[u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES]> {
let payload = json.get("applicationKeyAgreement")?.as_object()?;
if payload.get("algorithm").and_then(|value| value.as_str())
!= Some(crate::key_agreement::KEY_AGREEMENT_ALGORITHM)
{
return None;
}
let public_key = payload.get("publicKey")?.as_str()?.trim();
if public_key.is_empty() {
return None;
}
let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
.decode(public_key)
.ok()?;
if bytes.len() != crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES {
return None;
}
let mut out = [0u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES];
out.copy_from_slice(&bytes);
Some(out)
}
fn parse_transport_trust_device_id(json: &serde_json::Value) -> Option<String> {
json.get("transportTrust")
.and_then(|value| value.as_object())
.and_then(|payload| payload.get("deviceId"))
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "lowercase")]
pub enum NativeMessageAction {
Request,
Response,
Event,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct NativeMainMessage {
pub channel: String,
pub action: NativeMessageAction,
#[serde(skip_serializing_if = "Option::is_none")]
pub request_id: Option<String>,
pub payload: serde_json::Value,
pub timestamp: i64,
#[serde(skip_serializing_if = "Option::is_none")]
pub from: Option<String>,
}
pub const SESSION_TOKEN_CHANNEL: &str = "session-token";
impl NativeMainMessage {
pub fn is_handshake(&self) -> bool {
self.channel == "handshake"
}
pub fn is_session_token_presentation(&self) -> bool {
self.channel == SESSION_TOKEN_CHANNEL
}
pub fn is_session_token_response(&self) -> bool {
self.channel == SESSION_TOKEN_CHANNEL
&& matches!(self.action, NativeMessageAction::Response)
}
pub fn claimed_device_id(&self) -> Option<String> {
self.payload
.get("deviceId")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
pub fn presented_session_token(&self) -> Option<String> {
self.payload
.get("token")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
pub fn presented_session_token_payload(&self) -> Option<String> {
self.payload
.get("tokenPayload")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
pub fn session_token_approved(&self) -> Option<bool> {
self.payload
.get("approved")
.and_then(|value| value.as_bool())
}
pub fn approved_session_scope(&self) -> Option<String> {
self.payload
.get("scope")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
pub fn session_token_error(&self) -> Option<String> {
self.payload
.get("reason")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned)
}
pub fn session_token_presentation(token: &str) -> Self {
Self::session_token_presentation_with_payload(token, None)
}
pub fn session_token_presentation_with_payload(
token: &str,
token_payload: Option<&str>,
) -> Self {
Self::session_token_presentation_with_payload_and_device_id(token, token_payload, None)
}
pub fn session_token_presentation_with_payload_and_device_id(
token: &str,
token_payload: Option<&str>,
device_id: Option<&str>,
) -> Self {
#[cfg(target_arch = "wasm32")]
let timestamp = js_sys::Date::now() as i64;
#[cfg(not(target_arch = "wasm32"))]
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
let mut payload = serde_json::json!({ "token": token });
if let Some(token_payload) = token_payload {
if let Some(object) = payload.as_object_mut() {
object.insert(
"tokenPayload".to_string(),
serde_json::Value::String(token_payload.to_string()),
);
}
}
if let Some(device_id) = device_id.map(str::trim).filter(|value| !value.is_empty()) {
if let Some(object) = payload.as_object_mut() {
object.insert(
"deviceId".to_string(),
serde_json::Value::String(device_id.to_string()),
);
}
}
Self {
channel: SESSION_TOKEN_CHANNEL.to_string(),
action: NativeMessageAction::Request,
request_id: None,
payload,
timestamp,
from: None,
}
}
pub fn session_token_approval(scope: Option<&str>, connection_id: &str) -> Self {
#[cfg(target_arch = "wasm32")]
let timestamp = js_sys::Date::now() as i64;
#[cfg(not(target_arch = "wasm32"))]
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
Self {
channel: SESSION_TOKEN_CHANNEL.to_string(),
action: NativeMessageAction::Response,
request_id: None,
payload: serde_json::json!({
"approved": true,
"scope": scope,
"connectionId": connection_id,
}),
timestamp,
from: None,
}
}
pub fn session_token_rejection(reason: &str, connection_id: &str) -> Self {
#[cfg(target_arch = "wasm32")]
let timestamp = js_sys::Date::now() as i64;
#[cfg(not(target_arch = "wasm32"))]
let timestamp = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as i64)
.unwrap_or(0);
Self {
channel: SESSION_TOKEN_CHANNEL.to_string(),
action: NativeMessageAction::Response,
request_id: None,
payload: serde_json::json!({
"approved": false,
"reason": reason,
"connectionId": connection_id,
}),
timestamp,
from: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TypeScriptHandshakeCapabilities {
pub webrtc: Option<bool>,
pub moq: Option<bool>,
pub application_key_agreement: Option<bool>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TypeScriptHandshake {
pub action: Option<String>,
pub session_token: Option<String>,
pub session_token_payload: Option<String>,
pub claimed_device_id: Option<String>,
pub capabilities: Option<TypeScriptHandshakeCapabilities>,
pub application_key_agreement_public_key:
Option<[u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES]>,
}
#[derive(Debug, Clone)]
pub enum ParsedMainFrame {
NativeMessage(NativeMainMessage),
TypeScriptHandshake(TypeScriptHandshake),
TypeScriptJson(serde_json::Value),
Opaque,
}
#[derive(Debug, Clone)]
pub enum InspectedMainFrame {
NativeMessage {
message: NativeMainMessage,
handshake: Option<NativeHandshakeBinding>,
},
ForwardOpaque,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct NativeHandshakeBinding {
pub known_device_id: Option<String>,
pub claimed_device_id: Option<String>,
pub admitted_device_id: Option<String>,
pub authoritative_device_id_hint: Option<String>,
}
pub fn parse_main_frame(frame: &[u8]) -> ParsedMainFrame {
if let Ok(message) = serde_json::from_slice::<NativeMainMessage>(frame) {
return ParsedMainFrame::NativeMessage(message);
}
if frame.len() >= 2 && frame[0] == 0x00 {
if let Ok(json) = serde_json::from_slice::<serde_json::Value>(&frame[1..]) {
if json.get("type").and_then(|value| value.as_str()) == Some("handshake") {
return ParsedMainFrame::TypeScriptHandshake(TypeScriptHandshake {
action: json
.get("action")
.and_then(|value| value.as_str())
.map(ToOwned::to_owned),
session_token: json
.get("sessionToken")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned),
session_token_payload: json
.get("sessionTokenPayload")
.and_then(|value| value.as_str())
.map(str::trim)
.filter(|value| !value.is_empty())
.map(ToOwned::to_owned),
claimed_device_id: parse_transport_trust_device_id(&json),
capabilities: json
.get("capabilities")
.and_then(|value| value.as_object())
.map(|value| TypeScriptHandshakeCapabilities {
webrtc: value.get("webrtc").and_then(|flag| flag.as_bool()),
moq: value.get("moq").and_then(|flag| flag.as_bool()),
application_key_agreement: value
.get("applicationKeyAgreement")
.and_then(|flag| flag.as_bool()),
}),
application_key_agreement_public_key:
parse_application_key_agreement_public_key(&json),
});
}
return ParsedMainFrame::TypeScriptJson(json);
}
}
ParsedMainFrame::Opaque
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_native_handshake_message() {
let frame = serde_json::json!({
"channel": "handshake",
"action": "request",
"payload": { "deviceId": "device-1" },
"timestamp": 123,
});
match parse_main_frame(frame.to_string().as_bytes()) {
ParsedMainFrame::NativeMessage(message) => {
assert!(message.is_handshake());
assert_eq!(message.claimed_device_id().as_deref(), Some("device-1"));
}
other => panic!("expected native message, got {:?}", other),
}
}
#[test]
fn parses_typescript_handshake_message() {
let json = serde_json::json!({
"type": "handshake",
"action": "hello",
"sessionToken": "token-1",
"sessionTokenPayload": "payload-1",
"transportTrust": {
"version": 1,
"deviceId": "device-1",
"identityFingerprint": "fingerprint-1"
},
"capabilities": {
"webrtc": true,
"moq": false,
},
});
let mut frame = vec![0x00];
frame.extend_from_slice(json.to_string().as_bytes());
match parse_main_frame(&frame) {
ParsedMainFrame::TypeScriptHandshake(handshake) => {
assert_eq!(handshake.action.as_deref(), Some("hello"));
assert_eq!(handshake.session_token.as_deref(), Some("token-1"));
assert_eq!(
handshake.session_token_payload.as_deref(),
Some("payload-1")
);
assert_eq!(handshake.claimed_device_id.as_deref(), Some("device-1"));
assert_eq!(
handshake.capabilities.and_then(|value| value.webrtc),
Some(true)
);
}
other => panic!("expected typescript handshake, got {:?}", other),
}
}
#[test]
fn falls_back_to_opaque_for_unrecognized_bytes() {
assert!(matches!(
parse_main_frame(b"\x01\x02\x03"),
ParsedMainFrame::Opaque
));
}
#[test]
fn parses_pluto_signal_envelope_with_ts_prefix() {
let json = serde_json::json!({
"type": "#pluto-signal",
"content": {
"transport": "webrtc",
"type": "sdp",
"sdp": { "type": "answer", "sdp": "v=0\r\n" }
}
});
let mut frame = vec![0x00];
frame.extend_from_slice(json.to_string().as_bytes());
match parse_main_frame(&frame) {
ParsedMainFrame::TypeScriptJson(parsed) => {
assert_eq!(
parsed.get("type").and_then(|v| v.as_str()),
Some("#pluto-signal")
);
let content = parsed.get("content").unwrap();
assert_eq!(
content.get("transport").and_then(|v| v.as_str()),
Some("webrtc")
);
}
other => panic!("expected TypeScriptJson for #pluto-signal, got {:?}", other),
}
}
#[test]
fn raw_pluto_signal_without_prefix_is_opaque() {
let json = serde_json::json!({
"type": "#pluto-signal",
"content": { "transport": "webrtc", "type": "sdp" }
});
let frame = json.to_string().into_bytes();
match parse_main_frame(&frame) {
ParsedMainFrame::NativeMessage(_) => {
panic!("raw #pluto-signal should not parse as NativeMainMessage")
}
ParsedMainFrame::TypeScriptJson(_) => {
panic!("raw #pluto-signal without 0x00 prefix should not be TypeScriptJson")
}
ParsedMainFrame::Opaque => {}
other => panic!("expected Opaque, got {:?}", other),
}
}
#[test]
fn session_token_response_helpers_round_trip() {
let presentation = NativeMainMessage::session_token_presentation_with_payload(
"token-1",
Some("payload-1"),
);
assert!(presentation.is_session_token_presentation());
assert_eq!(
presentation.presented_session_token().as_deref(),
Some("token-1")
);
assert_eq!(
presentation.presented_session_token_payload().as_deref(),
Some("payload-1")
);
let presentation_with_device =
NativeMainMessage::session_token_presentation_with_payload_and_device_id(
"token-1",
Some("payload-1"),
Some("device-1"),
);
assert_eq!(
presentation_with_device.claimed_device_id().as_deref(),
Some("device-1")
);
let approved = NativeMainMessage::session_token_approval(Some("magic-link"), "conn-1");
assert!(approved.is_session_token_response());
assert_eq!(approved.session_token_approved(), Some(true));
assert_eq!(
approved.approved_session_scope().as_deref(),
Some("magic-link")
);
let rejected = NativeMainMessage::session_token_rejection("invalid-token", "conn-1");
assert!(rejected.is_session_token_response());
assert_eq!(rejected.session_token_approved(), Some(false));
assert_eq!(
rejected.session_token_error().as_deref(),
Some("invalid-token")
);
}
}