1use o_sfu_model::RecordingOptions;
2use serde::{Serialize, de::DeserializeOwned};
3use serde_json::Value;
4
5use super::{
6 AuthPayload, ClientBroadcastPayload, Envelope, PeerInfoPayload, PeerLeftPayload,
7 RecordingActionResult, RequestId, ServerBroadcastPayload, SessionDescriptionPayload,
8 SourceDescriptor, StreamIntentPayload, SubscribePayload, TrackBinding, WelcomePayload,
9};
10use crate::shared::{RecordingStateUpdate, UserInfo};
11
12const AUTH: &str = "auth";
13const BROADCAST: &str = "broadcast";
14const INFO: &str = "info";
15const OFFER: &str = "offer";
16const PEER_INFO: &str = "peerinfo";
17const PEER_JOINED: &str = "peerjoined";
18const PEER_LEFT: &str = "peerleft";
19const PUBLISH: &str = "publish";
20const RECORDING_CHANGE: &str = "recordingchange";
21const RENEGOTIATE: &str = "renegotiate";
22const START_RECORDING: &str = "startrecording";
23const STOP_RECORDING: &str = "stoprecording";
24const SUBSCRIBE: &str = "subscribe";
25const SOURCES: &str = "sources";
26const TRACKS: &str = "tracks";
27const UNPUBLISH: &str = "unpublish";
28const WELCOME: &str = "welcome";
29
30#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)]
31pub enum EnvelopeDecodeError {
32 #[error("unknown envelope tag: {0}")]
33 UnknownTag(String),
34 #[error("invalid payload for envelope tag: {0}")]
35 InvalidPayload(String),
36 #[error("unexpected payload for envelope tag: {0}")]
37 UnexpectedPayload(String),
38}
39
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub enum ClientMessage {
42 Auth(AuthPayload),
43 Publish(StreamIntentPayload),
44 Unpublish(StreamIntentPayload),
45 Subscribe(SubscribePayload),
46 Info(UserInfo),
47 Broadcast(ClientBroadcastPayload),
48}
49
50impl ClientMessage {
51 pub(crate) fn into_envelope(self) -> Result<Envelope, serde_json::Error> {
52 match self {
53 Self::Auth(payload) => encode_auth(payload),
54 Self::Publish(payload) => encode_message(PUBLISH, payload),
55 Self::Unpublish(payload) => encode_message(UNPUBLISH, payload),
56 Self::Subscribe(payload) => encode_message(SUBSCRIBE, payload),
57 Self::Info(payload) => encode_message(INFO, payload),
58 Self::Broadcast(payload) => encode_message(BROADCAST, payload),
59 }
60 }
61
62 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
63 match tag {
64 AUTH => parse_payload(tag, payload).map(Self::Auth),
65 PUBLISH => parse_payload(tag, payload).map(Self::Publish),
66 UNPUBLISH => parse_payload(tag, payload).map(Self::Unpublish),
67 SUBSCRIBE => parse_payload(tag, payload).map(Self::Subscribe),
68 INFO => parse_payload(tag, payload).map(Self::Info),
69 BROADCAST => parse_payload(tag, payload).map(Self::Broadcast),
70 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
71 }
72 }
73}
74
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub enum ClientRequest {
77 StartRecording(RecordingOptions),
78 StopRecording,
79}
80
81impl ClientRequest {
82 pub(crate) fn into_envelope(
83 self,
84 request_id: RequestId,
85 ) -> Result<Envelope, serde_json::Error> {
86 match self {
87 Self::StartRecording(payload) => encode_request(START_RECORDING, request_id, payload),
88 Self::StopRecording => Ok(Envelope::request(STOP_RECORDING, request_id, None)),
89 }
90 }
91
92 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
93 match tag {
94 START_RECORDING => parse_payload(tag, payload).map(Self::StartRecording),
95 STOP_RECORDING => {
96 ensure_empty_payload(tag, payload.as_ref())?;
97 Ok(Self::StopRecording)
98 }
99 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
100 }
101 }
102}
103
104#[derive(Debug, Clone, PartialEq, Eq)]
105pub enum ServerRequest {
106 Offer(SessionDescriptionPayload),
107 Renegotiate(SessionDescriptionPayload),
108}
109
110impl ServerRequest {
111 pub fn into_envelope(self, request_id: RequestId) -> Result<Envelope, serde_json::Error> {
118 match self {
119 Self::Offer(payload) => encode_request(OFFER, request_id, payload),
120 Self::Renegotiate(payload) => encode_request(RENEGOTIATE, request_id, payload),
121 }
122 }
123
124 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
125 match tag {
126 OFFER => parse_payload(tag, payload).map(Self::Offer),
127 RENEGOTIATE => parse_payload(tag, payload).map(Self::Renegotiate),
128 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
129 }
130 }
131}
132
133#[derive(Debug, Clone, PartialEq, Eq)]
134pub enum ClientResponse {
135 Offer(SessionDescriptionPayload),
136 Renegotiate(SessionDescriptionPayload),
137}
138
139impl ClientResponse {
140 pub(crate) fn into_envelope(
141 self,
142 response_to: RequestId,
143 ) -> Result<Envelope, serde_json::Error> {
144 match self {
145 Self::Offer(payload) => encode_response(OFFER, response_to, payload),
146 Self::Renegotiate(payload) => encode_response(RENEGOTIATE, response_to, payload),
147 }
148 }
149
150 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
151 match tag {
152 OFFER => parse_payload(tag, payload).map(Self::Offer),
153 RENEGOTIATE => parse_payload(tag, payload).map(Self::Renegotiate),
154 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
155 }
156 }
157}
158
159#[derive(Debug, Clone, PartialEq, Eq)]
160pub enum ServerMessage {
161 Welcome(WelcomePayload),
162 Tracks(Vec<TrackBinding>),
163 Sources(Vec<SourceDescriptor>),
164 PeerInfo(PeerInfoPayload),
165 PeerJoined(PeerInfoPayload),
166 PeerLeft(PeerLeftPayload),
167 Broadcast(ServerBroadcastPayload),
168 RecordingChange(RecordingStateUpdate),
169}
170
171impl ServerMessage {
172 pub fn into_envelope(self) -> Result<Envelope, serde_json::Error> {
179 match self {
180 Self::Welcome(payload) => encode_message(WELCOME, payload),
181 Self::Tracks(payload) => encode_message(TRACKS, payload),
182 Self::Sources(payload) => encode_message(SOURCES, payload),
183 Self::PeerInfo(payload) => encode_message(PEER_INFO, payload),
184 Self::PeerJoined(payload) => encode_message(PEER_JOINED, payload),
185 Self::PeerLeft(payload) => encode_message(PEER_LEFT, payload),
186 Self::Broadcast(payload) => encode_message(BROADCAST, payload),
187 Self::RecordingChange(payload) => encode_message(RECORDING_CHANGE, payload),
188 }
189 }
190
191 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
192 match tag {
193 WELCOME => parse_payload(tag, payload).map(Self::Welcome),
194 TRACKS => parse_payload(tag, payload).map(Self::Tracks),
195 SOURCES => parse_payload(tag, payload).map(Self::Sources),
196 PEER_INFO => parse_payload(tag, payload).map(Self::PeerInfo),
197 PEER_JOINED => parse_payload(tag, payload).map(Self::PeerJoined),
198 PEER_LEFT => parse_payload(tag, payload).map(Self::PeerLeft),
199 BROADCAST => parse_payload(tag, payload).map(Self::Broadcast),
200 RECORDING_CHANGE => parse_payload(tag, payload).map(Self::RecordingChange),
201 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
202 }
203 }
204}
205
206#[derive(Debug, Clone, PartialEq, Eq)]
207pub enum ServerResponse {
208 StartRecording(RecordingActionResult),
209 StopRecording(RecordingActionResult),
210}
211
212impl ServerResponse {
213 pub fn into_envelope(self, response_to: RequestId) -> Result<Envelope, serde_json::Error> {
220 match self {
221 Self::StartRecording(payload) => encode_response(START_RECORDING, response_to, payload),
222 Self::StopRecording(payload) => encode_response(STOP_RECORDING, response_to, payload),
223 }
224 }
225
226 pub(crate) fn decode(tag: &str, payload: Option<Value>) -> Result<Self, EnvelopeDecodeError> {
227 match tag {
228 START_RECORDING => parse_payload(tag, payload).map(Self::StartRecording),
229 STOP_RECORDING => parse_payload(tag, payload).map(Self::StopRecording),
230 _ => Err(EnvelopeDecodeError::UnknownTag(tag.to_owned())),
231 }
232 }
233}
234
235#[cfg(feature = "client")]
236fn encode_auth(payload: AuthPayload) -> Result<Envelope, serde_json::Error> {
237 encode_message(AUTH, payload)
238}
239
240#[cfg(not(feature = "client"))]
241#[expect(
242 clippy::unreachable,
243 reason = "defensive guard: this build has no Serialize impl for AuthPayload, so nothing can construct a ClientMessage::Auth to encode"
244)]
245fn encode_auth(_payload: AuthPayload) -> Result<Envelope, serde_json::Error> {
250 unreachable!(
251 "ClientMessage::Auth is only ever encoded by a client-role host (the `client` feature)"
252 )
253}
254
255fn encode_message<T: Serialize>(tag: &str, payload: T) -> Result<Envelope, serde_json::Error> {
256 Ok(Envelope::message(tag, Some(serde_json::to_value(payload)?)))
257}
258
259fn encode_request<T: Serialize>(
260 tag: &str,
261 request_id: RequestId,
262 payload: T,
263) -> Result<Envelope, serde_json::Error> {
264 Ok(Envelope::request(
265 tag,
266 request_id,
267 Some(serde_json::to_value(payload)?),
268 ))
269}
270
271fn encode_response<T: Serialize>(
272 tag: &str,
273 response_to: RequestId,
274 payload: T,
275) -> Result<Envelope, serde_json::Error> {
276 Ok(Envelope::response(
277 tag,
278 response_to,
279 Some(serde_json::to_value(payload)?),
280 ))
281}
282
283fn parse_payload<T: DeserializeOwned>(
284 tag: &str,
285 payload: Option<Value>,
286) -> Result<T, EnvelopeDecodeError> {
287 serde_json::from_value(
288 payload.ok_or_else(|| EnvelopeDecodeError::InvalidPayload(tag.to_owned()))?,
289 )
290 .map_err(|_error| EnvelopeDecodeError::InvalidPayload(tag.to_owned()))
291}
292
293fn ensure_empty_payload(tag: &str, payload: Option<&Value>) -> Result<(), EnvelopeDecodeError> {
294 if payload.is_some() {
295 return Err(EnvelopeDecodeError::UnexpectedPayload(tag.to_owned()));
296 }
297 Ok(())
298}