Skip to main content

openrtc/
native_protocol.rs

1use base64::Engine as _;
2use serde::{Deserialize, Serialize};
3
4fn parse_application_key_agreement_public_key(
5    json: &serde_json::Value,
6) -> Option<[u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES]> {
7    let payload = json.get("applicationKeyAgreement")?.as_object()?;
8    if payload.get("algorithm").and_then(|value| value.as_str())
9        != Some(crate::key_agreement::KEY_AGREEMENT_ALGORITHM)
10    {
11        return None;
12    }
13    let public_key = payload.get("publicKey")?.as_str()?.trim();
14    if public_key.is_empty() {
15        return None;
16    }
17    let bytes = base64::engine::general_purpose::URL_SAFE_NO_PAD
18        .decode(public_key)
19        .ok()?;
20    if bytes.len() != crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES {
21        return None;
22    }
23    let mut out = [0u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES];
24    out.copy_from_slice(&bytes);
25    Some(out)
26}
27
28fn parse_transport_trust_device_id(json: &serde_json::Value) -> Option<String> {
29    json.get("transportTrust")
30        .and_then(|value| value.as_object())
31        .and_then(|payload| payload.get("deviceId"))
32        .and_then(|value| value.as_str())
33        .map(str::trim)
34        .filter(|value| !value.is_empty())
35        .map(ToOwned::to_owned)
36}
37
38#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
39#[serde(rename_all = "lowercase")]
40pub enum NativeMessageAction {
41    Request,
42    Response,
43    Event,
44}
45
46#[derive(Debug, Clone, Serialize, Deserialize)]
47#[serde(rename_all = "camelCase")]
48pub struct NativeMainMessage {
49    pub channel: String,
50    pub action: NativeMessageAction,
51    #[serde(skip_serializing_if = "Option::is_none")]
52    pub request_id: Option<String>,
53    pub payload: serde_json::Value,
54    pub timestamp: i64,
55    #[serde(skip_serializing_if = "Option::is_none")]
56    pub from: Option<String>,
57}
58
59/// Channel name used for SDK-owned session token presentation.
60/// The connecting client sends this as the first native message after
61/// connection when a compound ticket contained an embedded token.
62pub const SESSION_TOKEN_CHANNEL: &str = "session-token";
63pub const PEER_DATA_CHANNEL: &str = "peer-data";
64/// `send_peer` is the bounded, ordered message surface. Larger payloads use a
65/// named stream or a protocol package so lifecycle/control frames cannot be
66/// stalled behind bulk transfer data on the native-main stream.
67pub const MAX_PEER_MESSAGE_BYTES: usize = 64 * 1024;
68
69/// Lifetime contract for the stream carrying a session-token presentation.
70///
71/// This describes the stream, not the sender runtime or grant scope. Browser
72/// bridges commonly use a bounded request/response exchange, while native
73/// peers retain the admitted stream as their generation-bound control route.
74#[derive(Debug, Clone, Copy, Default, Serialize, Deserialize, PartialEq, Eq)]
75#[serde(rename_all = "kebab-case")]
76pub enum SessionTokenStreamContract {
77    #[default]
78    OneShotAdmission,
79    PersistentControl,
80}
81
82/// Reverse-direction admission carried in the response to a typed reciprocal
83/// request. Keeping it on the request's stream binds both proofs to the same
84/// physical generation and avoids selecting a stale endpoint connection after
85/// network replacement.
86#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
87#[serde(rename_all = "camelCase")]
88pub struct ReciprocalSessionTokenPresentation {
89    pub presentation_id: String,
90    pub token: String,
91    pub token_payload: String,
92    pub device_id: String,
93    pub stream_contract: SessionTokenStreamContract,
94}
95
96#[derive(Debug, Clone, PartialEq, Eq)]
97pub enum ReciprocalSessionTokenPresentationState {
98    Absent,
99    Malformed,
100    Present(ReciprocalSessionTokenPresentation),
101}
102
103impl NativeMainMessage {
104    pub fn is_handshake(&self) -> bool {
105        self.channel == "handshake"
106    }
107
108    pub fn is_session_token_presentation(&self) -> bool {
109        self.channel == SESSION_TOKEN_CHANNEL && matches!(self.action, NativeMessageAction::Request)
110    }
111
112    pub fn is_session_token_response(&self) -> bool {
113        self.channel == SESSION_TOKEN_CHANNEL
114            && matches!(self.action, NativeMessageAction::Response)
115    }
116
117    pub fn is_session_token_response_ack(&self, connection_id: &str) -> bool {
118        self.channel == SESSION_TOKEN_CHANNEL
119            && matches!(self.action, NativeMessageAction::Event)
120            && self
121                .payload
122                .get("responseAck")
123                .and_then(|value| value.as_bool())
124                == Some(true)
125            && self
126                .payload
127                .get("connectionId")
128                .and_then(|value| value.as_str())
129                == Some(connection_id)
130    }
131
132    pub fn is_peer_data(&self) -> bool {
133        self.channel == PEER_DATA_CHANNEL
134    }
135
136    pub fn peer_data_payload(&self) -> Option<Vec<u8>> {
137        let payload_base64 = self
138            .payload
139            .get("payloadBase64")
140            .and_then(|value| value.as_str())
141            .map(str::trim)
142            .filter(|value| !value.is_empty())?;
143        base64::engine::general_purpose::STANDARD
144            .decode(payload_base64.as_bytes())
145            .ok()
146    }
147
148    pub fn claimed_device_id(&self) -> Option<String> {
149        self.payload
150            .get("deviceId")
151            .and_then(|value| value.as_str())
152            .map(str::trim)
153            .filter(|value| !value.is_empty())
154            .map(ToOwned::to_owned)
155    }
156
157    pub fn presented_session_token(&self) -> Option<String> {
158        self.payload
159            .get("token")
160            .and_then(|value| value.as_str())
161            .map(str::trim)
162            .filter(|value| !value.is_empty())
163            .map(ToOwned::to_owned)
164    }
165
166    pub fn presented_session_token_payload(&self) -> Option<String> {
167        self.payload
168            .get("tokenPayload")
169            .and_then(|value| value.as_str())
170            .map(str::trim)
171            .filter(|value| !value.is_empty())
172            .map(ToOwned::to_owned)
173    }
174
175    pub fn session_token_stream_contract(&self) -> SessionTokenStreamContract {
176        self.payload
177            .get("streamContract")
178            .cloned()
179            .and_then(|value| serde_json::from_value(value).ok())
180            .unwrap_or_default()
181    }
182
183    /// Whether the presenter has lost the reciprocal inbound proof for this
184    /// physical generation and is asking the peer's admission actor to
185    /// re-present its own token. Missing on older peers means `false`.
186    pub fn session_token_reciprocal_requested(&self) -> bool {
187        self.payload
188            .get("reciprocalRequested")
189            .and_then(|value| value.as_bool())
190            .unwrap_or(false)
191    }
192
193    pub fn reciprocal_session_token_mode(&self) -> Option<&str> {
194        self.payload
195            .get("reciprocalMode")
196            .and_then(|value| value.as_str())
197    }
198
199    pub fn reciprocal_session_token_presentation_state(
200        &self,
201    ) -> ReciprocalSessionTokenPresentationState {
202        let Some(presentation) = self.payload.get("reciprocalPresentation").cloned() else {
203            return ReciprocalSessionTokenPresentationState::Absent;
204        };
205        let Ok(parsed) = serde_json::from_value::<ReciprocalSessionTokenPresentation>(presentation)
206        else {
207            return ReciprocalSessionTokenPresentationState::Malformed;
208        };
209        if parsed.presentation_id.trim().is_empty()
210            || parsed.token.trim().is_empty()
211            || parsed.token_payload.trim().is_empty()
212            || parsed.device_id.trim().is_empty()
213        {
214            return ReciprocalSessionTokenPresentationState::Malformed;
215        }
216        ReciprocalSessionTokenPresentationState::Present(parsed)
217    }
218
219    pub fn reciprocal_session_token_presentation_id(&self) -> Option<String> {
220        self.payload
221            .get("reciprocalPresentation")
222            .and_then(|value| value.get("presentationId"))
223            .and_then(|value| value.as_str())
224            .map(str::trim)
225            .filter(|value| !value.is_empty())
226            .map(ToOwned::to_owned)
227    }
228
229    /// New native peers opt into an explicit approval-delivery ACK before a
230    /// persistent control stream participates in canonical stream arbitration.
231    pub fn session_token_response_ack_requested(&self) -> bool {
232        self.payload
233            .get("responseAckRequested")
234            .and_then(|value| value.as_bool())
235            .unwrap_or(false)
236    }
237
238    pub fn with_session_token_stream_contract(
239        mut self,
240        contract: SessionTokenStreamContract,
241    ) -> Self {
242        if let Some(payload) = self.payload.as_object_mut() {
243            payload.insert(
244                "streamContract".to_string(),
245                serde_json::to_value(contract)
246                    .expect("session-token stream contract must serialize"),
247            );
248        }
249        self
250    }
251
252    pub fn with_session_token_reciprocal_requested(mut self, requested: bool) -> Self {
253        if requested {
254            if let Some(payload) = self.payload.as_object_mut() {
255                payload.insert(
256                    "reciprocalRequested".to_string(),
257                    serde_json::Value::Bool(true),
258                );
259                payload.insert(
260                    "reciprocalMode".to_string(),
261                    serde_json::Value::String("inline-v1".to_string()),
262                );
263            }
264        }
265        self
266    }
267
268    pub fn session_token_approved(&self) -> Option<bool> {
269        self.payload
270            .get("approved")
271            .and_then(|value| value.as_bool())
272    }
273
274    pub fn approved_session_scope(&self) -> Option<String> {
275        self.payload
276            .get("scope")
277            .and_then(|value| value.as_str())
278            .map(str::trim)
279            .filter(|value| !value.is_empty())
280            .map(ToOwned::to_owned)
281    }
282
283    pub fn session_token_error(&self) -> Option<String> {
284        self.payload
285            .get("reason")
286            .and_then(|value| value.as_str())
287            .map(str::trim)
288            .filter(|value| !value.is_empty())
289            .map(ToOwned::to_owned)
290    }
291
292    /// Build a session-token presentation message for the connecting client
293    /// to send to the host.
294    pub fn session_token_presentation(token: &str) -> Self {
295        Self::session_token_presentation_with_payload(token, None)
296    }
297
298    pub fn session_token_presentation_with_payload(
299        token: &str,
300        token_payload: Option<&str>,
301    ) -> Self {
302        Self::session_token_presentation_with_payload_and_device_id(token, token_payload, None)
303    }
304
305    pub fn session_token_presentation_with_payload_and_device_id(
306        token: &str,
307        token_payload: Option<&str>,
308        device_id: Option<&str>,
309    ) -> Self {
310        #[cfg(target_arch = "wasm32")]
311        let timestamp = js_sys::Date::now() as i64;
312        #[cfg(not(target_arch = "wasm32"))]
313        let timestamp = std::time::SystemTime::now()
314            .duration_since(std::time::UNIX_EPOCH)
315            .map(|d| d.as_millis() as i64)
316            .unwrap_or(0);
317
318        let mut payload = serde_json::json!({
319            "token": token,
320            "streamContract": SessionTokenStreamContract::OneShotAdmission,
321            "responseAckRequested": true,
322        });
323        if let Some(token_payload) = token_payload {
324            if let Some(object) = payload.as_object_mut() {
325                object.insert(
326                    "tokenPayload".to_string(),
327                    serde_json::Value::String(token_payload.to_string()),
328                );
329            }
330        }
331        if let Some(device_id) = device_id.map(str::trim).filter(|value| !value.is_empty()) {
332            if let Some(object) = payload.as_object_mut() {
333                object.insert(
334                    "deviceId".to_string(),
335                    serde_json::Value::String(device_id.to_string()),
336                );
337            }
338        }
339
340        Self {
341            channel: SESSION_TOKEN_CHANNEL.to_string(),
342            action: NativeMessageAction::Request,
343            request_id: None,
344            payload,
345            timestamp,
346            from: None,
347        }
348    }
349
350    pub fn session_token_approval(scope: Option<&str>, connection_id: &str) -> Self {
351        #[cfg(target_arch = "wasm32")]
352        let timestamp = js_sys::Date::now() as i64;
353        #[cfg(not(target_arch = "wasm32"))]
354        let timestamp = std::time::SystemTime::now()
355            .duration_since(std::time::UNIX_EPOCH)
356            .map(|d| d.as_millis() as i64)
357            .unwrap_or(0);
358
359        Self {
360            channel: SESSION_TOKEN_CHANNEL.to_string(),
361            action: NativeMessageAction::Response,
362            request_id: None,
363            payload: serde_json::json!({
364                "approved": true,
365                "scope": scope,
366                "connectionId": connection_id,
367            }),
368            timestamp,
369            from: None,
370        }
371    }
372
373    pub fn with_reciprocal_session_token_presentation(
374        mut self,
375        presentation: ReciprocalSessionTokenPresentation,
376    ) -> Self {
377        if let Some(payload) = self.payload.as_object_mut() {
378            payload.insert(
379                "reciprocalMode".to_string(),
380                serde_json::Value::String("inline-v1".to_string()),
381            );
382            payload.insert(
383                "reciprocalPresentation".to_string(),
384                serde_json::to_value(presentation)
385                    .expect("reciprocal session-token presentation must serialize"),
386            );
387        }
388        self
389    }
390
391    pub fn session_token_rejection(reason: &str, connection_id: &str) -> Self {
392        #[cfg(target_arch = "wasm32")]
393        let timestamp = js_sys::Date::now() as i64;
394        #[cfg(not(target_arch = "wasm32"))]
395        let timestamp = std::time::SystemTime::now()
396            .duration_since(std::time::UNIX_EPOCH)
397            .map(|d| d.as_millis() as i64)
398            .unwrap_or(0);
399
400        Self {
401            channel: SESSION_TOKEN_CHANNEL.to_string(),
402            action: NativeMessageAction::Response,
403            request_id: None,
404            payload: serde_json::json!({
405                "approved": false,
406                "reason": reason,
407                "connectionId": connection_id,
408            }),
409            timestamp,
410            from: None,
411        }
412    }
413
414    pub fn session_token_response_ack(connection_id: &str) -> Self {
415        #[cfg(target_arch = "wasm32")]
416        let timestamp = js_sys::Date::now() as i64;
417        #[cfg(not(target_arch = "wasm32"))]
418        let timestamp = std::time::SystemTime::now()
419            .duration_since(std::time::UNIX_EPOCH)
420            .map(|d| d.as_millis() as i64)
421            .unwrap_or(0);
422
423        Self {
424            channel: SESSION_TOKEN_CHANNEL.to_string(),
425            action: NativeMessageAction::Event,
426            request_id: None,
427            payload: serde_json::json!({
428                "responseAck": true,
429                "connectionId": connection_id,
430            }),
431            timestamp,
432            from: None,
433        }
434    }
435
436    pub fn with_reciprocal_session_token_ack(
437        mut self,
438        presentation_id: &str,
439        accepted: bool,
440        scope: Option<&str>,
441        reason: Option<&str>,
442    ) -> Self {
443        if let Some(payload) = self.payload.as_object_mut() {
444            payload.insert(
445                "reciprocalPresentationId".to_string(),
446                serde_json::Value::String(presentation_id.to_string()),
447            );
448            payload.insert(
449                "reciprocalAccepted".to_string(),
450                serde_json::Value::Bool(accepted),
451            );
452            if let Some(scope) = scope {
453                payload.insert(
454                    "reciprocalScope".to_string(),
455                    serde_json::Value::String(scope.to_string()),
456                );
457            }
458            if let Some(reason) = reason {
459                payload.insert(
460                    "reciprocalReason".to_string(),
461                    serde_json::Value::String(reason.to_string()),
462                );
463            }
464        }
465        self
466    }
467
468    pub fn reciprocal_session_token_ack_presentation_id(&self) -> Option<String> {
469        self.payload
470            .get("reciprocalPresentationId")
471            .and_then(|value| value.as_str())
472            .map(str::trim)
473            .filter(|value| !value.is_empty())
474            .map(ToOwned::to_owned)
475    }
476
477    pub fn reciprocal_session_token_ack_accepted(&self) -> Option<bool> {
478        self.payload
479            .get("reciprocalAccepted")
480            .and_then(|value| value.as_bool())
481    }
482
483    pub fn reciprocal_session_token_ack_scope(&self) -> Option<String> {
484        self.payload
485            .get("reciprocalScope")
486            .and_then(|value| value.as_str())
487            .map(str::trim)
488            .filter(|value| !value.is_empty())
489            .map(ToOwned::to_owned)
490    }
491
492    pub fn peer_data(payload: &[u8]) -> Self {
493        #[cfg(target_arch = "wasm32")]
494        let timestamp = js_sys::Date::now() as i64;
495        #[cfg(not(target_arch = "wasm32"))]
496        let timestamp = std::time::SystemTime::now()
497            .duration_since(std::time::UNIX_EPOCH)
498            .map(|d| d.as_millis() as i64)
499            .unwrap_or(0);
500
501        Self {
502            channel: PEER_DATA_CHANNEL.to_string(),
503            action: NativeMessageAction::Event,
504            request_id: None,
505            payload: serde_json::json!({
506                "payloadBase64": base64::engine::general_purpose::STANDARD.encode(payload),
507            }),
508            timestamp,
509            from: None,
510        }
511    }
512}
513
514#[derive(Debug, Clone, PartialEq, Eq)]
515pub struct TypeScriptHandshakeCapabilities {
516    pub webrtc: Option<bool>,
517    pub moq: Option<bool>,
518    pub ble: Option<bool>,
519    pub application_key_agreement: Option<bool>,
520    pub scoped_webrtc_signal_stream: Option<bool>,
521}
522
523#[derive(Debug, Clone, PartialEq, Eq)]
524pub struct TypeScriptHandshake {
525    pub action: Option<String>,
526    pub session_token: Option<String>,
527    pub session_token_payload: Option<String>,
528    pub claimed_device_id: Option<String>,
529    pub capabilities: Option<TypeScriptHandshakeCapabilities>,
530    pub application_key_agreement_public_key:
531        Option<[u8; crate::key_agreement::KEY_AGREEMENT_PUBLIC_KEY_BYTES]>,
532}
533
534#[derive(Debug, Clone)]
535pub enum ParsedMainFrame {
536    NativeMessage(NativeMainMessage),
537    TypeScriptHandshake(TypeScriptHandshake),
538    TypeScriptJson(serde_json::Value),
539    Opaque,
540}
541
542#[derive(Debug, Clone)]
543pub enum InspectedMainFrame {
544    NativeMessage {
545        message: NativeMainMessage,
546        handshake: Option<NativeHandshakeBinding>,
547    },
548    /// SDK-owned response for a session-token presentation received on an
549    /// application-forwarded one-shot main stream. The host adapter writes
550    /// this response verbatim, then asks the client to run post-admission
551    /// lifecycle work only after the wire flush completes.
552    SessionTokenResponse {
553        response: NativeMainMessage,
554        accepted: bool,
555        is_first_presentation: bool,
556    },
557    ForwardOpaque,
558}
559
560#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
561#[serde(rename_all = "camelCase")]
562pub struct NativeHandshakeBinding {
563    pub known_device_id: Option<String>,
564    pub claimed_device_id: Option<String>,
565    pub admitted_device_id: Option<String>,
566    pub authoritative_device_id_hint: Option<String>,
567}
568
569pub fn parse_main_frame(frame: &[u8]) -> ParsedMainFrame {
570    if let Ok(message) = serde_json::from_slice::<NativeMainMessage>(frame) {
571        return ParsedMainFrame::NativeMessage(message);
572    }
573
574    if frame.len() >= 2 && frame[0] == 0x00 {
575        if let Ok(json) = serde_json::from_slice::<serde_json::Value>(&frame[1..]) {
576            if json.get("type").and_then(|value| value.as_str()) == Some("handshake") {
577                return ParsedMainFrame::TypeScriptHandshake(TypeScriptHandshake {
578                    action: json
579                        .get("action")
580                        .and_then(|value| value.as_str())
581                        .map(ToOwned::to_owned),
582                    session_token: json
583                        .get("sessionToken")
584                        .and_then(|value| value.as_str())
585                        .map(str::trim)
586                        .filter(|value| !value.is_empty())
587                        .map(ToOwned::to_owned),
588                    session_token_payload: json
589                        .get("sessionTokenPayload")
590                        .and_then(|value| value.as_str())
591                        .map(str::trim)
592                        .filter(|value| !value.is_empty())
593                        .map(ToOwned::to_owned),
594                    claimed_device_id: parse_transport_trust_device_id(&json),
595                    capabilities: json
596                        .get("capabilities")
597                        .and_then(|value| value.as_object())
598                        .map(|value| TypeScriptHandshakeCapabilities {
599                            webrtc: value.get("webrtc").and_then(|flag| flag.as_bool()),
600                            moq: value.get("moq").and_then(|flag| flag.as_bool()),
601                            ble: value.get("ble").and_then(|flag| flag.as_bool()),
602                            application_key_agreement: value
603                                .get("applicationKeyAgreement")
604                                .and_then(|flag| flag.as_bool()),
605                            scoped_webrtc_signal_stream: value
606                                .get("scopedWebRtcSignalStream")
607                                .and_then(|flag| flag.as_bool()),
608                        }),
609                    application_key_agreement_public_key:
610                        parse_application_key_agreement_public_key(&json),
611                });
612            }
613            return ParsedMainFrame::TypeScriptJson(json);
614        }
615    }
616
617    ParsedMainFrame::Opaque
618}
619
620#[cfg(test)]
621mod tests {
622    use super::*;
623
624    #[test]
625    fn parses_native_handshake_message() {
626        let frame = serde_json::json!({
627            "channel": "handshake",
628            "action": "request",
629            "payload": { "deviceId": "device-1" },
630            "timestamp": 123,
631        });
632
633        match parse_main_frame(frame.to_string().as_bytes()) {
634            ParsedMainFrame::NativeMessage(message) => {
635                assert!(message.is_handshake());
636                assert_eq!(message.claimed_device_id().as_deref(), Some("device-1"));
637            }
638            other => panic!("expected native message, got {:?}", other),
639        }
640    }
641
642    #[test]
643    fn parses_typescript_handshake_message() {
644        let json = serde_json::json!({
645            "type": "handshake",
646            "action": "hello",
647            "sessionToken": "token-1",
648            "sessionTokenPayload": "payload-1",
649            "transportTrust": {
650                "version": 1,
651                "deviceId": "device-1",
652                "identityFingerprint": "fingerprint-1"
653            },
654            "capabilities": {
655                "webrtc": true,
656                "moq": false,
657                "ble": true,
658                "scopedWebRtcSignalStream": true,
659            },
660        });
661        let mut frame = vec![0x00];
662        frame.extend_from_slice(json.to_string().as_bytes());
663
664        match parse_main_frame(&frame) {
665            ParsedMainFrame::TypeScriptHandshake(handshake) => {
666                assert_eq!(handshake.action.as_deref(), Some("hello"));
667                assert_eq!(handshake.session_token.as_deref(), Some("token-1"));
668                assert_eq!(
669                    handshake.session_token_payload.as_deref(),
670                    Some("payload-1")
671                );
672                assert_eq!(handshake.claimed_device_id.as_deref(), Some("device-1"));
673                let capabilities = handshake.capabilities.expect("capabilities");
674                assert_eq!(capabilities.webrtc, Some(true));
675                assert_eq!(capabilities.moq, Some(false));
676                assert_eq!(capabilities.ble, Some(true));
677                assert_eq!(capabilities.scoped_webrtc_signal_stream, Some(true));
678            }
679            other => panic!("expected typescript handshake, got {:?}", other),
680        }
681    }
682
683    #[test]
684    fn falls_back_to_opaque_for_unrecognized_bytes() {
685        assert!(matches!(
686            parse_main_frame(b"\x01\x02\x03"),
687            ParsedMainFrame::Opaque
688        ));
689    }
690
691    #[test]
692    fn parses_peer_data_native_message() {
693        let message = NativeMainMessage::peer_data(b"hello peer");
694        let frame = serde_json::to_vec(&message).expect("peer data should serialize");
695
696        match parse_main_frame(&frame) {
697            ParsedMainFrame::NativeMessage(parsed) => {
698                assert!(parsed.is_peer_data());
699                assert_eq!(
700                    parsed.peer_data_payload().as_deref(),
701                    Some(b"hello peer".as_slice())
702                );
703            }
704            other => panic!("expected native peer data message, got {:?}", other),
705        }
706    }
707
708    #[test]
709    fn session_token_stream_contract_round_trips_and_legacy_defaults_to_one_shot() {
710        let persistent = NativeMainMessage::session_token_presentation("token")
711            .with_session_token_stream_contract(SessionTokenStreamContract::PersistentControl);
712        let encoded = serde_json::to_vec(&persistent).expect("serialize persistent contract");
713        let decoded: NativeMainMessage =
714            serde_json::from_slice(&encoded).expect("deserialize persistent contract");
715        assert_eq!(
716            decoded.session_token_stream_contract(),
717            SessionTokenStreamContract::PersistentControl,
718        );
719
720        let mut legacy = NativeMainMessage::session_token_presentation("legacy-token");
721        legacy
722            .payload
723            .as_object_mut()
724            .expect("presentation payload")
725            .remove("streamContract");
726        assert_eq!(
727            legacy.session_token_stream_contract(),
728            SessionTokenStreamContract::OneShotAdmission,
729        );
730        assert!(!legacy.session_token_reciprocal_requested());
731    }
732
733    #[test]
734    fn session_token_reciprocal_request_round_trips_and_defaults_false() {
735        let legacy = NativeMainMessage::session_token_presentation("legacy-token");
736        assert!(!legacy.session_token_reciprocal_requested());
737        assert_eq!(legacy.reciprocal_session_token_mode(), None);
738
739        let requested = NativeMainMessage::session_token_presentation("token")
740            .with_session_token_reciprocal_requested(true);
741        let encoded = serde_json::to_vec(&requested).expect("serialize reciprocal request");
742        let decoded: NativeMainMessage =
743            serde_json::from_slice(&encoded).expect("deserialize reciprocal request");
744        assert!(decoded.session_token_reciprocal_requested());
745        assert_eq!(decoded.reciprocal_session_token_mode(), Some("inline-v1"));
746    }
747
748    #[test]
749    fn reciprocal_session_token_presentation_round_trips_on_approval() {
750        let response = NativeMainMessage::session_token_approval(Some("user-device"), "conn-1")
751            .with_reciprocal_session_token_presentation(ReciprocalSessionTokenPresentation {
752                presentation_id: "presentation-1".to_string(),
753                token: "reverse-token".to_string(),
754                token_payload: "reverse-payload".to_string(),
755                device_id: "device-1".to_string(),
756                stream_contract: SessionTokenStreamContract::PersistentControl,
757            });
758        let encoded = serde_json::to_vec(&response).expect("serialize reciprocal response");
759        let decoded: NativeMainMessage =
760            serde_json::from_slice(&encoded).expect("deserialize reciprocal response");
761        assert_eq!(decoded.reciprocal_session_token_mode(), Some("inline-v1"));
762        assert_eq!(
763            decoded.reciprocal_session_token_presentation_state(),
764            ReciprocalSessionTokenPresentationState::Present(ReciprocalSessionTokenPresentation {
765                presentation_id: "presentation-1".to_string(),
766                token: "reverse-token".to_string(),
767                token_payload: "reverse-payload".to_string(),
768                device_id: "device-1".to_string(),
769                stream_contract: SessionTokenStreamContract::PersistentControl,
770            })
771        );
772
773        assert_eq!(
774            NativeMainMessage::session_token_approval(None, "conn-1")
775                .reciprocal_session_token_presentation_state(),
776            ReciprocalSessionTokenPresentationState::Absent,
777        );
778
779        let malformed = NativeMainMessage::session_token_approval(None, "conn-1")
780            .with_reciprocal_session_token_presentation(ReciprocalSessionTokenPresentation {
781                presentation_id: "presentation-2".to_string(),
782                token: "".to_string(),
783                token_payload: "reverse-payload".to_string(),
784                device_id: "device-1".to_string(),
785                stream_contract: SessionTokenStreamContract::OneShotAdmission,
786            });
787        assert_eq!(
788            malformed.reciprocal_session_token_presentation_state(),
789            ReciprocalSessionTokenPresentationState::Malformed,
790        );
791    }
792
793    #[test]
794    fn parses_pluto_signal_envelope_with_ts_prefix() {
795        let json = serde_json::json!({
796            "type": "#pluto-signal",
797            "content": {
798                "transport": "webrtc",
799                "type": "sdp",
800                "sdp": { "type": "answer", "sdp": "v=0\r\n" }
801            }
802        });
803        let mut frame = vec![0x00];
804        frame.extend_from_slice(json.to_string().as_bytes());
805
806        match parse_main_frame(&frame) {
807            ParsedMainFrame::TypeScriptJson(parsed) => {
808                assert_eq!(
809                    parsed.get("type").and_then(|v| v.as_str()),
810                    Some("#pluto-signal")
811                );
812                let content = parsed.get("content").unwrap();
813                assert_eq!(
814                    content.get("transport").and_then(|v| v.as_str()),
815                    Some("webrtc")
816                );
817            }
818            other => panic!("expected TypeScriptJson for #pluto-signal, got {:?}", other),
819        }
820    }
821
822    #[test]
823    fn raw_pluto_signal_without_prefix_is_opaque() {
824        let json = serde_json::json!({
825            "type": "#pluto-signal",
826            "content": { "transport": "webrtc", "type": "sdp" }
827        });
828        let frame = json.to_string().into_bytes();
829
830        match parse_main_frame(&frame) {
831            ParsedMainFrame::NativeMessage(_) => {
832                panic!("raw #pluto-signal should not parse as NativeMainMessage")
833            }
834            ParsedMainFrame::TypeScriptJson(_) => {
835                panic!("raw #pluto-signal without 0x00 prefix should not be TypeScriptJson")
836            }
837            ParsedMainFrame::Opaque => {}
838            other => panic!("expected Opaque, got {:?}", other),
839        }
840    }
841
842    #[test]
843    fn session_token_response_helpers_round_trip() {
844        let presentation = NativeMainMessage::session_token_presentation_with_payload(
845            "token-1",
846            Some("payload-1"),
847        );
848        assert!(presentation.is_session_token_presentation());
849        assert_eq!(
850            presentation.presented_session_token().as_deref(),
851            Some("token-1")
852        );
853        assert_eq!(
854            presentation.presented_session_token_payload().as_deref(),
855            Some("payload-1")
856        );
857        assert!(presentation.session_token_response_ack_requested());
858
859        let presentation_with_device =
860            NativeMainMessage::session_token_presentation_with_payload_and_device_id(
861                "token-1",
862                Some("payload-1"),
863                Some("device-1"),
864            );
865        assert_eq!(
866            presentation_with_device.claimed_device_id().as_deref(),
867            Some("device-1")
868        );
869
870        let approved = NativeMainMessage::session_token_approval(Some("magic-link"), "conn-1");
871        assert!(approved.is_session_token_response());
872        assert_eq!(approved.session_token_approved(), Some(true));
873        assert_eq!(
874            approved.approved_session_scope().as_deref(),
875            Some("magic-link")
876        );
877
878        let rejected = NativeMainMessage::session_token_rejection("invalid-token", "conn-1");
879        assert!(rejected.is_session_token_response());
880        assert_eq!(rejected.session_token_approved(), Some(false));
881        assert_eq!(
882            rejected.session_token_error().as_deref(),
883            Some("invalid-token")
884        );
885
886        let ack = NativeMainMessage::session_token_response_ack("conn-1");
887        assert!(!ack.is_session_token_presentation());
888        assert!(ack.is_session_token_response_ack("conn-1"));
889        assert!(!ack.is_session_token_response_ack("conn-2"));
890
891        let reciprocal_ack = NativeMainMessage::session_token_response_ack("conn-1")
892            .with_reciprocal_session_token_ack("presentation-1", true, Some("user-device"), None);
893        assert_eq!(
894            reciprocal_ack.payload.get("reciprocalPresentationId"),
895            Some(&serde_json::Value::String("presentation-1".to_string())),
896        );
897        assert_eq!(
898            reciprocal_ack.payload.get("reciprocalAccepted"),
899            Some(&serde_json::Value::Bool(true)),
900        );
901    }
902}