diff --git a/rust/crates/truapi-codegen/tests/golden/dispatcher.rs b/rust/crates/truapi-codegen/tests/golden/dispatcher.rs index 25f2fb13..f5629d77 100644 --- a/rust/crates/truapi-codegen/tests/golden/dispatcher.rs +++ b/rust/crates/truapi-codegen/tests/golden/dispatcher.rs @@ -15,6 +15,7 @@ use truapi::api::{ Entropy, LocalStorage, Notifications, + P2pMedia, Payment, Permissions, Preimage, @@ -46,6 +47,7 @@ where register_entropy(dispatcher, host.clone()); register_local_storage(dispatcher, host.clone()); register_notifications(dispatcher, host.clone()); + register_p2p_media(dispatcher, host.clone()); register_payment(dispatcher, host.clone()); register_permissions(dispatcher, host.clone()); register_preimage(dispatcher, host.clone()); @@ -1182,6 +1184,238 @@ where } } +fn register_p2p_media

(dispatcher: &mut Dispatcher, host: Arc

) +where + P: P2pMedia + Send + Sync + 'static, +{ + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_STATUS, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pStatusRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pStatusResponse = match host.status(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_ROOM_CREATE, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomCreateRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomCreateResponse = match host.room_create(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_ROOM_JOIN, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomJoinRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomJoinResponse = match host.room_join(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_ROOM_LEAVE, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomLeaveRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomLeaveResponse = match host.room_leave(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_ENDPOINT_REFRESH, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pEndpointRefreshRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pEndpointRefreshResponse = match host.endpoint_refresh(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_PUBLISH, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pPublishRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pPublishResponse = match host.publish(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_UNPUBLISH, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pUnpublishRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pUnpublishResponse = match host.unpublish(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host; + dispatcher.on_subscription(wire_table::P_2_P_MEDIA_ROOM_EVENTS, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomEventsRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Err(encode_versioned_interrupt_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let stream = match host.room_events(&cx, request).await { + Ok(sub) => sub, + Err(err) => { + return Err(encode_versioned_interrupt_payload(err, target_version)); + } + }; + Ok(subscription_stream::(stream)) + }) + }); + } +} + fn register_payment

(dispatcher: &mut Dispatcher, host: Arc

) where P: Payment + Send + Sync + 'static, diff --git a/rust/crates/truapi-codegen/tests/golden/wire_table.rs b/rust/crates/truapi-codegen/tests/golden/wire_table.rs index 12a72a94..40bea985 100644 --- a/rust/crates/truapi-codegen/tests/golden/wire_table.rs +++ b/rust/crates/truapi-codegen/tests/golden/wire_table.rs @@ -460,6 +460,56 @@ pub const COIN_PAYMENT_LISTEN_FOR_PAYMENT: SubscriptionFrameIds = SubscriptionFr receive_id: 163, }; +/// Wire discriminants for `p2p_media_status`. +pub const P_2_P_MEDIA_STATUS: RequestFrameIds = RequestFrameIds { + request_id: 164, + response_id: 165, +}; + +/// Wire discriminants for `p2p_media_room_create`. +pub const P_2_P_MEDIA_ROOM_CREATE: RequestFrameIds = RequestFrameIds { + request_id: 166, + response_id: 167, +}; + +/// Wire discriminants for `p2p_media_room_join`. +pub const P_2_P_MEDIA_ROOM_JOIN: RequestFrameIds = RequestFrameIds { + request_id: 168, + response_id: 169, +}; + +/// Wire discriminants for `p2p_media_room_leave`. +pub const P_2_P_MEDIA_ROOM_LEAVE: RequestFrameIds = RequestFrameIds { + request_id: 170, + response_id: 171, +}; + +/// Wire discriminants for `p2p_media_endpoint_refresh`. +pub const P_2_P_MEDIA_ENDPOINT_REFRESH: RequestFrameIds = RequestFrameIds { + request_id: 172, + response_id: 173, +}; + +/// Wire discriminants for `p2p_media_publish`. +pub const P_2_P_MEDIA_PUBLISH: RequestFrameIds = RequestFrameIds { + request_id: 174, + response_id: 175, +}; + +/// Wire discriminants for `p2p_media_unpublish`. +pub const P_2_P_MEDIA_UNPUBLISH: RequestFrameIds = RequestFrameIds { + request_id: 176, + response_id: 177, +}; + +/// Wire discriminants for `p2p_media_room_events`. +pub const P_2_P_MEDIA_ROOM_EVENTS: SubscriptionFrameIds = SubscriptionFrameIds { + start_id: 178, + stop_id: 179, + interrupt_id: 180, + receive_id: 181, +}; + /// The full wire table. Ordering is part of the wire protocol; /// only ever append. Removed methods leave their slot empty. pub const WIRE_TABLE: &[WireEntry] = &[ @@ -719,4 +769,36 @@ pub const WIRE_TABLE: &[WireEntry] = &[ method: "coin_payment_listen_for_payment", kind: WireKind::Subscription(COIN_PAYMENT_LISTEN_FOR_PAYMENT), }, + WireEntry { + method: "p2p_media_status", + kind: WireKind::Request(P_2_P_MEDIA_STATUS), + }, + WireEntry { + method: "p2p_media_room_create", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_CREATE), + }, + WireEntry { + method: "p2p_media_room_join", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_JOIN), + }, + WireEntry { + method: "p2p_media_room_leave", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_LEAVE), + }, + WireEntry { + method: "p2p_media_endpoint_refresh", + kind: WireKind::Request(P_2_P_MEDIA_ENDPOINT_REFRESH), + }, + WireEntry { + method: "p2p_media_publish", + kind: WireKind::Request(P_2_P_MEDIA_PUBLISH), + }, + WireEntry { + method: "p2p_media_unpublish", + kind: WireKind::Request(P_2_P_MEDIA_UNPUBLISH), + }, + WireEntry { + method: "p2p_media_room_events", + kind: WireKind::Subscription(P_2_P_MEDIA_ROOM_EVENTS), + }, ]; diff --git a/rust/crates/truapi-server/src/generated/dispatcher.rs b/rust/crates/truapi-server/src/generated/dispatcher.rs index 231dd362..e811958b 100644 --- a/rust/crates/truapi-server/src/generated/dispatcher.rs +++ b/rust/crates/truapi-server/src/generated/dispatcher.rs @@ -8,8 +8,8 @@ use parity_scale_codec::Decode; use truapi::CallContext; use truapi::api::{ - Account, Chain, Chat, CoinPayment, Entropy, LocalStorage, Notifications, Payment, Permissions, - Preimage, ResourceAllocation, Signing, StatementStore, System, Theme, + Account, Chain, Chat, CoinPayment, Entropy, LocalStorage, Notifications, P2pMedia, Payment, + Permissions, Preimage, ResourceAllocation, Signing, StatementStore, System, Theme, }; use truapi::versioned::{self, Versioned}; @@ -33,6 +33,7 @@ where register_entropy(dispatcher, host.clone()); register_local_storage(dispatcher, host.clone()); register_notifications(dispatcher, host.clone()); + register_p2p_media(dispatcher, host.clone()); register_payment(dispatcher, host.clone()); register_permissions(dispatcher, host.clone()); register_preimage(dispatcher, host.clone()); @@ -1345,6 +1346,294 @@ where } } +fn register_p2p_media

(dispatcher: &mut Dispatcher, host: Arc

) +where + P: P2pMedia + Send + Sync + 'static, +{ + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_STATUS, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pStatusRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pStatusError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pStatusResponse = + match host.status(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_ROOM_CREATE, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomCreateRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pRoomCreateError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomCreateResponse = + match host.room_create(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_ROOM_JOIN, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomJoinRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pRoomJoinError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomJoinResponse = + match host.room_join(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_ROOM_LEAVE, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomLeaveRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pRoomLeaveError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pRoomLeaveResponse = + match host.room_leave(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host.clone(); + dispatcher.on_request(wire_table::P_2_P_MEDIA_ENDPOINT_REFRESH, move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pEndpointRefreshRequest = match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError = + truapi::CallError::MalformedFrame { reason: err.to_string() }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pEndpointRefreshResponse = match host.endpoint_refresh(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }); + } + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_PUBLISH, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pPublishRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pPublishError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pPublishResponse = + match host.publish(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host.clone(); + dispatcher.on_request( + wire_table::P_2_P_MEDIA_UNPUBLISH, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pUnpublishRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pUnpublishError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Ok(encode_versioned_err_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let response: versioned::p2p_media::HostP2pUnpublishResponse = + match host.unpublish(&cx, request).await { + Ok(value) => value, + Err(err) => { + return Ok(encode_versioned_err_payload(err, target_version)); + } + }; + Ok(encode_versioned_ok_payload(response)) + }) + }, + ); + } + { + let host = host; + dispatcher.on_subscription( + wire_table::P_2_P_MEDIA_ROOM_EVENTS, + move |request_id: String, bytes: Vec| { + let host = host.clone(); + Box::pin(async move { + let request: versioned::p2p_media::HostP2pRoomEventsRequest = + match Decode::decode(&mut &bytes[..]) { + Ok(request) => request, + Err(err) => { + let error: truapi::CallError< + versioned::p2p_media::HostP2pRoomEventsError, + > = truapi::CallError::MalformedFrame { + reason: err.to_string(), + }; + return Err(encode_versioned_interrupt_payload( + error, + ::LATEST, + )); + } + }; + let target_version = request.version(); + let cx = CallContext::with_request_id(request_id.clone()); + let stream = match host.room_events(&cx, request).await { + Ok(sub) => sub, + Err(err) => { + return Err(encode_versioned_interrupt_payload(err, target_version)); + } + }; + Ok(subscription_stream::< + versioned::p2p_media::HostP2pRoomEventsItem, + _, + >(stream)) + }) + }, + ); + } +} + fn register_payment

(dispatcher: &mut Dispatcher, host: Arc

) where P: Payment + Send + Sync + 'static, diff --git a/rust/crates/truapi-server/src/generated/wire_table.rs b/rust/crates/truapi-server/src/generated/wire_table.rs index 12a72a94..40bea985 100644 --- a/rust/crates/truapi-server/src/generated/wire_table.rs +++ b/rust/crates/truapi-server/src/generated/wire_table.rs @@ -460,6 +460,56 @@ pub const COIN_PAYMENT_LISTEN_FOR_PAYMENT: SubscriptionFrameIds = SubscriptionFr receive_id: 163, }; +/// Wire discriminants for `p2p_media_status`. +pub const P_2_P_MEDIA_STATUS: RequestFrameIds = RequestFrameIds { + request_id: 164, + response_id: 165, +}; + +/// Wire discriminants for `p2p_media_room_create`. +pub const P_2_P_MEDIA_ROOM_CREATE: RequestFrameIds = RequestFrameIds { + request_id: 166, + response_id: 167, +}; + +/// Wire discriminants for `p2p_media_room_join`. +pub const P_2_P_MEDIA_ROOM_JOIN: RequestFrameIds = RequestFrameIds { + request_id: 168, + response_id: 169, +}; + +/// Wire discriminants for `p2p_media_room_leave`. +pub const P_2_P_MEDIA_ROOM_LEAVE: RequestFrameIds = RequestFrameIds { + request_id: 170, + response_id: 171, +}; + +/// Wire discriminants for `p2p_media_endpoint_refresh`. +pub const P_2_P_MEDIA_ENDPOINT_REFRESH: RequestFrameIds = RequestFrameIds { + request_id: 172, + response_id: 173, +}; + +/// Wire discriminants for `p2p_media_publish`. +pub const P_2_P_MEDIA_PUBLISH: RequestFrameIds = RequestFrameIds { + request_id: 174, + response_id: 175, +}; + +/// Wire discriminants for `p2p_media_unpublish`. +pub const P_2_P_MEDIA_UNPUBLISH: RequestFrameIds = RequestFrameIds { + request_id: 176, + response_id: 177, +}; + +/// Wire discriminants for `p2p_media_room_events`. +pub const P_2_P_MEDIA_ROOM_EVENTS: SubscriptionFrameIds = SubscriptionFrameIds { + start_id: 178, + stop_id: 179, + interrupt_id: 180, + receive_id: 181, +}; + /// The full wire table. Ordering is part of the wire protocol; /// only ever append. Removed methods leave their slot empty. pub const WIRE_TABLE: &[WireEntry] = &[ @@ -719,4 +769,36 @@ pub const WIRE_TABLE: &[WireEntry] = &[ method: "coin_payment_listen_for_payment", kind: WireKind::Subscription(COIN_PAYMENT_LISTEN_FOR_PAYMENT), }, + WireEntry { + method: "p2p_media_status", + kind: WireKind::Request(P_2_P_MEDIA_STATUS), + }, + WireEntry { + method: "p2p_media_room_create", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_CREATE), + }, + WireEntry { + method: "p2p_media_room_join", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_JOIN), + }, + WireEntry { + method: "p2p_media_room_leave", + kind: WireKind::Request(P_2_P_MEDIA_ROOM_LEAVE), + }, + WireEntry { + method: "p2p_media_endpoint_refresh", + kind: WireKind::Request(P_2_P_MEDIA_ENDPOINT_REFRESH), + }, + WireEntry { + method: "p2p_media_publish", + kind: WireKind::Request(P_2_P_MEDIA_PUBLISH), + }, + WireEntry { + method: "p2p_media_unpublish", + kind: WireKind::Request(P_2_P_MEDIA_UNPUBLISH), + }, + WireEntry { + method: "p2p_media_room_events", + kind: WireKind::Subscription(P_2_P_MEDIA_ROOM_EVENTS), + }, ]; diff --git a/rust/crates/truapi-server/src/runtime.rs b/rust/crates/truapi-server/src/runtime.rs index ee990b2d..8e39140a 100644 --- a/rust/crates/truapi-server/src/runtime.rs +++ b/rust/crates/truapi-server/src/runtime.rs @@ -59,8 +59,8 @@ use futures::{FutureExt, StreamExt, pin_mut}; use parity_scale_codec::Encode; use tracing::{info, instrument, warn}; use truapi::api::{ - Account, Chain, Chat, CoinPayment, Entropy, LocalStorage, Notifications, Payment, Permissions, - Preimage, ResourceAllocation, Signing, System, Theme, + Account, Chain, Chat, CoinPayment, Entropy, LocalStorage, Notifications, P2pMedia, Payment, + Permissions, Preimage, ResourceAllocation, Signing, System, Theme, }; use truapi::v01; use truapi::versioned::account::{ @@ -1750,6 +1750,7 @@ const PAYMENTS_NOT_IMPLEMENTED: &str = "Payments are not supported in dot.li"; impl Chat for ProductRuntimeHost {} impl CoinPayment for ProductRuntimeHost {} +impl P2pMedia for ProductRuntimeHost {} impl Payment for ProductRuntimeHost { #[instrument(skip_all, fields(runtime.method = "payment.balance_subscribe"))] async fn balance_subscribe( diff --git a/rust/crates/truapi/src/api.rs b/rust/crates/truapi/src/api.rs index 957509e4..f96d9613 100644 --- a/rust/crates/truapi/src/api.rs +++ b/rust/crates/truapi/src/api.rs @@ -7,6 +7,7 @@ pub mod coin_payment; pub mod entropy; pub mod local_storage; pub mod notifications; +pub mod p2p_media; pub mod payment; pub mod permissions; pub mod preimage; @@ -23,6 +24,7 @@ pub use coin_payment::CoinPayment; pub use entropy::Entropy; pub use local_storage::LocalStorage; pub use notifications::Notifications; +pub use p2p_media::P2pMedia; pub use payment::Payment; pub use permissions::Permissions; pub use preimage::Preimage; @@ -41,6 +43,7 @@ pub trait TrUApi: + Entropy + LocalStorage + Notifications + + P2pMedia + Payment + Permissions + Preimage @@ -62,6 +65,7 @@ impl TrUApi for T where + Entropy + LocalStorage + Notifications + + P2pMedia + Payment + Permissions + Preimage diff --git a/rust/crates/truapi/src/api/p2p_media.rs b/rust/crates/truapi/src/api/p2p_media.rs new file mode 100644 index 00000000..e505cb38 --- /dev/null +++ b/rust/crates/truapi/src/api/p2p_media.rs @@ -0,0 +1,187 @@ +//! Unified [`P2pMedia`] trait. + +use crate::versioned::p2p_media::{ + HostP2pEndpointRefreshError, HostP2pEndpointRefreshRequest, HostP2pEndpointRefreshResponse, + HostP2pPublishError, HostP2pPublishRequest, HostP2pPublishResponse, HostP2pRoomCreateError, + HostP2pRoomCreateRequest, HostP2pRoomCreateResponse, HostP2pRoomEventsError, + HostP2pRoomEventsItem, HostP2pRoomEventsRequest, HostP2pRoomJoinError, HostP2pRoomJoinRequest, + HostP2pRoomJoinResponse, HostP2pRoomLeaveError, HostP2pRoomLeaveRequest, + HostP2pRoomLeaveResponse, HostP2pStatusError, HostP2pStatusRequest, HostP2pStatusResponse, + HostP2pUnpublishError, HostP2pUnpublishRequest, HostP2pUnpublishResponse, +}; +use crate::wire; +use crate::{CallContext, CallError, Subscription}; + +/// Peer-to-peer media rooms (MoQ-over-iroh). +/// +/// The host embeds a headless iroh/MoQ node and exposes rooms - behind +/// the RFC 0002 grant model. Products reach their room through a loopback moq +/// relay the host serves on `127.0.0.1`: publishes under `self/…` fan out to +/// peers, remote peers' broadcasts appear under `room//…`. +pub trait P2pMedia: Send + Sync { + /// Probe the host's p2p capability + node info. Never prompts. + /// + /// ```ts + /// const result = await truapi.p2PMedia.status(); + /// assert(result.isOk(), "status failed:", result); + /// console.log("p2p available:", result.value.available); + /// ``` + #[wire(request_id = 164)] + async fn status( + &self, + _cx: &CallContext, + _request: HostP2pStatusRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Create a room (host side). THE prompting method: resolves + /// `RemotePermission::MediaP2p` - even for receive-only rooms (peers learn + /// the user's network address) - plus `DevicePermission::{Camera, + /// Microphone}` per the requested directions, in ONE prompt, persisted per + /// RFC 0002. Grants (loopback connect allowance, inline playback/autoplay, + /// lifecycle) apply for the room's lifetime. + /// + /// ```ts + /// const result = await truapi.p2PMedia.roomCreate({ + /// directions: { + /// publishVideo: true, + /// publishAudio: true, + /// receiveVideo: true, + /// receiveAudio: true, + /// }, + /// purpose: "Video room", + /// displayName: undefined, + /// }); + /// assert(result.isOk(), "roomCreate failed:", result); + /// console.log("room ticket:", result.value.ticket); + /// ``` + #[wire(request_id = 166)] + async fn room_create( + &self, + _cx: &CallContext, + _request: HostP2pRoomCreateRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Join a room via its invite ticket string. Same gating and response + /// shape as `room_create`. + /// + /// ```ts + /// const result = await truapi.p2PMedia.roomJoin({ + /// ticket: "room-invite-ticket", + /// directions: { + /// publishVideo: false, + /// publishAudio: false, + /// receiveVideo: true, + /// receiveAudio: true, + /// }, + /// purpose: "Video room", + /// displayName: undefined, + /// }); + /// assert(result.isOk(), "roomJoin failed:", result); + /// console.log("joined room:", result.value.room); + /// ``` + #[wire(request_id = 168)] + async fn room_join( + &self, + _cx: &CallContext, + _request: HostP2pRoomJoinRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Leave a room and restore the default sandbox posture. De-escalation + /// only - no prompt (mirrors the rt-session `session_close`). + /// + /// ```ts + /// const result = await truapi.p2PMedia.roomLeave({ room: 1n }); + /// assert(result.isOk(), "roomLeave failed:", result); + /// ``` + #[wire(request_id = 170)] + async fn room_leave( + &self, + _cx: &CallContext, + _request: HostP2pRoomLeaveRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Re-issue the loopback endpoint (token/cert rotation). Valid room + /// required, no prompt - mirrors the rt-session `relay_token`. + /// + /// ```ts + /// const result = await truapi.p2PMedia.endpointRefresh({ room: 1n }); + /// assert(result.isOk(), "endpointRefresh failed:", result); + /// console.log("fresh endpoint:", result.value.endpoint.wtUrl); + /// ``` + #[wire(request_id = 172)] + async fn endpoint_refresh( + &self, + _cx: &CallContext, + _request: HostP2pEndpointRefreshRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Offer broadcast names (relative to the product's `self/` scope) to the + /// room. Valid room required; the product must be publishing them into the + /// loopback relay (or start within the host's patience). + /// + /// ```ts + /// const result = await truapi.p2PMedia.publish({ + /// room: 1n, + /// names: ["camera"], + /// }); + /// assert(result.isOk(), "publish failed:", result); + /// ``` + #[wire(request_id = 174)] + async fn publish( + &self, + _cx: &CallContext, + _request: HostP2pPublishRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Withdraw broadcast names from the room. + /// + /// ```ts + /// const result = await truapi.p2PMedia.unpublish({ + /// room: 1n, + /// names: ["camera"], + /// }); + /// assert(result.isOk(), "unpublish failed:", result); + /// ``` + #[wire(request_id = 176)] + async fn unpublish( + &self, + _cx: &CallContext, + _request: HostP2pUnpublishRequest, + ) -> Result> { + Err(CallError::unavailable()) + } + + /// Room membership + broadcast lifecycle + rt-session lifecycle events. + /// Holding this subscription is the KEEP-ALIVE signal: + /// it keeps the node running through backgrounding grace + /// windows, so a live media session must keep it open. + /// + /// ```ts + /// import { firstValueFrom, from } from "rxjs"; + /// + /// const event = await firstValueFrom( + /// from(truapi.p2PMedia.roomEvents({ room: 1n })), + /// ); + /// console.log("room event:", event); + /// ``` + #[wire(start_id = 178)] + async fn room_events( + &self, + _cx: &CallContext, + _request: HostP2pRoomEventsRequest, + ) -> Result, CallError> { + Err(CallError::unavailable()) + } +} diff --git a/rust/crates/truapi/src/v01.rs b/rust/crates/truapi/src/v01.rs index 8b34df5a..dc2185d6 100644 --- a/rust/crates/truapi/src/v01.rs +++ b/rust/crates/truapi/src/v01.rs @@ -8,6 +8,7 @@ mod common; mod entropy; mod local_storage; mod notifications; +mod p2p_media; mod payment; mod permissions; mod preimage; @@ -26,6 +27,7 @@ pub use common::*; pub use entropy::*; pub use local_storage::*; pub use notifications::*; +pub use p2p_media::*; pub use payment::*; pub use permissions::*; pub use preimage::*; diff --git a/rust/crates/truapi/src/v01/p2p_media.rs b/rust/crates/truapi/src/v01/p2p_media.rs new file mode 100644 index 00000000..f0952605 --- /dev/null +++ b/rust/crates/truapi/src/v01/p2p_media.rs @@ -0,0 +1,212 @@ +use parity_scale_codec::{Decode, Encode}; + +/// Opaque p2p room handle, host-minted, scoped to the issuing product instance. +pub type P2pRoomId = u64; + +/// Which media directions a room wants. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Encode, Decode)] +pub struct RtDirections { + /// Publish camera video into the room. + pub publish_video: bool, + /// Publish microphone audio into the room. + pub publish_audio: bool, + /// Receive remote peers' video. + pub receive_video: bool, + /// Receive remote peers' audio. + pub receive_audio: bool, +} + +/// The loopback endpoint a product dials with `@moq/*` / raw WebTransport +/// (`RtRelayConfig`-shaped, loopback edition). The token rides both URLs as +/// `?jwt=…` and scopes the session to the room's root: publish under `self/` +/// only, subscribe everything (bridged peers appear at `room//`). +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pEndpoint { + /// `https://127.0.0.1:/?jwt=` - WebTransport. + pub wt_url: String, + /// sha-256 (hex) of the relay's self-signed cert, for + /// `serverCertificateHashes`. + pub cert_sha256: String, + /// `ws://127.0.0.1:/?jwt=` - WebSocket fallback. + pub ws_url: String, + /// Token expiry, unix millis. Refresh via `endpoint_refresh`. + pub expires_at_ms: u64, +} + +/// Shared error union for every [`P2pMedia`](crate::api::P2pMedia) method. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub enum P2pError { + /// The user declined `MediaP2p` or a required device permission. + PermissionDenied, + /// The invite ticket failed to parse or has expired. + InvalidTicket, + /// All bootstrap peers offline / gossip join failed. + JoinFailed { + /// Human-readable failure diagnostic. + reason: String, + }, + /// The room handle does not name a live room of this product instance. + RoomNotFound, + /// A named broadcast never appeared in the loopback relay. + BroadcastMissing { + /// The broadcast name that never appeared. + name: String, + }, + /// The host's per-product concurrent room cap was reached. + TooManyRooms, + /// Caller modality may not perform this operation (e.g. capture + /// directions from a widget, or any room from a headless worker). + NotAllowedForModality, + /// This host/platform cannot run a p2p node. + Unsupported, + /// Catch-all. + Unknown { + /// Human-readable failure diagnostic. + reason: String, + }, +} + +/// `status` response: host p2p capability + node info. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pStatusResponse { + /// Whether this host can run p2p media rooms. + pub available: bool, + /// The node's stable iroh endpoint id (hex), when running. + pub endpoint_id: Option, + /// Number of live rooms held by the calling product. + pub num_rooms: u32, +} + +/// `room_create` request. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pRoomCreateRequest { + /// Requested media directions; publishing folds a camera/mic prompt. + pub directions: RtDirections, + /// Short human-readable purpose, shown in the permission prompt. + pub purpose: String, + /// Per-room presence display name. + pub display_name: Option, +} + +/// Shared response for `room_create` / `room_join`: the opaque handle, its +/// invite ticket, and the loopback endpoint. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pRoomResponse { + /// Opaque room handle - passed back to the other `P2pMedia` calls. + pub room: P2pRoomId, + /// The invite: a self-describing compact ticket string. + pub ticket: String, + /// The loopback endpoint the product dials. + pub endpoint: HostP2pEndpoint, +} + +/// `room_join` request. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pRoomJoinRequest { + /// The invite ticket string received out-of-band. + pub ticket: String, + /// Requested media directions; publishing folds a camera/mic prompt. + pub directions: RtDirections, + /// Short human-readable purpose, shown in the permission prompt. + pub purpose: String, + /// Per-room presence display name. + pub display_name: Option, +} + +/// `room_leave` request (unit response). +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pRoomLeaveRequest { + /// The room to leave. + pub room: P2pRoomId, +} + +/// `endpoint_refresh` request. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pEndpointRefreshRequest { + /// The room whose loopback endpoint to re-issue. + pub room: P2pRoomId, +} + +/// `endpoint_refresh` response. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pEndpointRefreshResponse { + /// The fresh loopback endpoint (rotated token/cert). + pub endpoint: HostP2pEndpoint, +} + +/// `publish` request (unit response). +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pPublishRequest { + /// The room to offer the broadcasts to. + pub room: P2pRoomId, + /// Broadcast names, relative to the product's `self/` scope. + pub names: Vec, +} + +/// `unpublish` request (unit response). +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pUnpublishRequest { + /// The room to withdraw the broadcasts from. + pub room: P2pRoomId, + /// Broadcast names, relative to the product's `self/` scope. + pub names: Vec, +} + +/// `room_events` subscription request. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub struct HostP2pRoomEventsRequest { + /// The room whose events to stream. + pub room: P2pRoomId, +} + +/// `room_events` subscription item: roster + broadcast lifecycle + +/// rt-session lifecycle events. +#[derive(Debug, Clone, PartialEq, Eq, Encode, Decode)] +pub enum HostP2pRoomEvent { + /// Emitted once on subscribe; grants are live. + Active, + /// A peer joined the room. + PeerJoined { + /// The peer's endpoint id (hex). + peer: String, + /// The peer's presence display name, if it announced one. + display_name: Option, + }, + /// A peer left the room. + PeerLeft { + /// The peer's endpoint id (hex). + peer: String, + }, + /// The peer's broadcast is now served by the local loopback relay at + /// `room//` under the product's scope root. + BroadcastAdded { + /// The publishing peer's endpoint id (hex). + peer: String, + /// The broadcast name. + name: String, + }, + /// The peer's broadcast disappeared from the loopback relay. + BroadcastRemoved { + /// The publishing peer's endpoint id (hex). + peer: String, + /// The broadcast name. + name: String, + }, + /// Loopback endpoint rotated (token/cert) - re-dial with the new config. + EndpointChanged { + /// The fresh loopback endpoint. + endpoint: HostP2pEndpoint, + }, + /// Platform is about to freeze the product. + Suspending { + /// Grace window before the freeze, in milliseconds. + grace_ms: u32, + }, + /// The product returned to the foreground. + Resumed, + /// Room ended by the host; default posture restored. + Revoked { + /// Human-readable revocation reason. + reason: String, + }, +} diff --git a/rust/crates/truapi/src/v01/permissions.rs b/rust/crates/truapi/src/v01/permissions.rs index c4873a6f..e07a5208 100644 --- a/rust/crates/truapi/src/v01/permissions.rs +++ b/rust/crates/truapi/src/v01/permissions.rs @@ -53,6 +53,11 @@ pub enum RemotePermission { /// Submitting statements on behalf of the user via `remote_statement_store_submit`. #[display("submit statements")] StatementSubmit, + /// Peer-to-peer media rooms: direct connections to other participants + /// (they learn this device's network address) plus the host's local + /// media relay, for the lifetime of a room. + #[display("peer-to-peer media rooms")] + MediaP2p, } /// remote-permission request (RFC 0002). diff --git a/rust/crates/truapi/src/versioned.rs b/rust/crates/truapi/src/versioned.rs index 9da72067..dc72c792 100644 --- a/rust/crates/truapi/src/versioned.rs +++ b/rust/crates/truapi/src/versioned.rs @@ -37,6 +37,7 @@ pub mod coin_payment; pub mod entropy; pub mod local_storage; pub mod notifications; +pub mod p2p_media; pub mod payment; pub mod permissions; pub mod preimage; diff --git a/rust/crates/truapi/src/versioned/p2p_media.rs b/rust/crates/truapi/src/versioned/p2p_media.rs new file mode 100644 index 00000000..929092fc --- /dev/null +++ b/rust/crates/truapi/src/versioned/p2p_media.rs @@ -0,0 +1,30 @@ +//! Versioned wrappers for [`P2pMedia`](crate::api::P2pMedia) methods. + +use crate::v01; + +truapi_macros::versioned_type! { + pub enum HostP2pStatusRequest { V1 } + pub enum HostP2pStatusResponse { V1 => v01::HostP2pStatusResponse } + pub enum HostP2pStatusError { V1 => v01::P2pError } + pub enum HostP2pRoomCreateRequest { V1 => v01::HostP2pRoomCreateRequest } + pub enum HostP2pRoomCreateResponse { V1 => v01::HostP2pRoomResponse } + pub enum HostP2pRoomCreateError { V1 => v01::P2pError } + pub enum HostP2pRoomJoinRequest { V1 => v01::HostP2pRoomJoinRequest } + pub enum HostP2pRoomJoinResponse { V1 => v01::HostP2pRoomResponse } + pub enum HostP2pRoomJoinError { V1 => v01::P2pError } + pub enum HostP2pRoomLeaveRequest { V1 => v01::HostP2pRoomLeaveRequest } + pub enum HostP2pRoomLeaveResponse { V1 } + pub enum HostP2pRoomLeaveError { V1 => v01::P2pError } + pub enum HostP2pEndpointRefreshRequest { V1 => v01::HostP2pEndpointRefreshRequest } + pub enum HostP2pEndpointRefreshResponse { V1 => v01::HostP2pEndpointRefreshResponse } + pub enum HostP2pEndpointRefreshError { V1 => v01::P2pError } + pub enum HostP2pPublishRequest { V1 => v01::HostP2pPublishRequest } + pub enum HostP2pPublishResponse { V1 } + pub enum HostP2pPublishError { V1 => v01::P2pError } + pub enum HostP2pUnpublishRequest { V1 => v01::HostP2pUnpublishRequest } + pub enum HostP2pUnpublishResponse { V1 } + pub enum HostP2pUnpublishError { V1 => v01::P2pError } + pub enum HostP2pRoomEventsRequest { V1 => v01::HostP2pRoomEventsRequest } + pub enum HostP2pRoomEventsItem { V1 => v01::HostP2pRoomEvent } + pub enum HostP2pRoomEventsError { V1 => v01::P2pError } +}