Skip to main content

o_sfu_protocol/core/
request_flow.rs

1use super::{
2    Command, Commands, FlushMode, NegotiationKind, NegotiationRejection, PendingRequest,
3    PendingRequestKind, ProtocolCore, REQUEST_TIMEOUT_MS, close_for_protocol_error,
4};
5use crate::signaling::{
6    ClientEnvelope, ClientRequest, ClientResponse, RecordingOptions, RequestId, ServerRequest,
7    ServerResponse, SessionDescriptionPayload,
8};
9
10impl ProtocolCore {
11    pub fn start_recording(&mut self, options: RecordingOptions) -> Vec<Command> {
12        begin_request(self, ClientRequest::StartRecording(options))
13    }
14
15    pub fn stop_recording(&mut self) -> Vec<Command> {
16        begin_request(self, ClientRequest::StopRecording)
17    }
18
19    /// Replies to the currently pending negotiation request.
20    ///
21    /// The host must echo the exact `request_id` and `kind` from
22    /// [`Command::ApplyNegotiation`]. Mismatches are ignored so a stale or
23    /// reordered SDP answer cannot accidentally resolve the wrong negotiation.
24    pub fn submit_negotiation_answer(
25        &mut self,
26        request_id: &RequestId,
27        kind: NegotiationKind,
28        sdp: impl Into<String>,
29    ) -> Vec<Command> {
30        if !self.phase.resolve_negotiation(request_id, kind) {
31            return Vec::new();
32        }
33        let payload = SessionDescriptionPayload {
34            sdp: sdp.into(),
35            upload_slots: Vec::new(),
36        };
37        let response = match kind {
38            NegotiationKind::Offer => ClientResponse::Offer(payload),
39            NegotiationKind::Renegotiate => ClientResponse::Renegotiate(payload),
40        };
41        let Some(envelope) = ClientEnvelope::Response {
42            response_to: request_id.clone(),
43            response,
44        }
45        .into_envelope()
46        .ok() else {
47            return Vec::new();
48        };
49        self.outbound_batch.enqueue(envelope, FlushMode::Immediate)
50    }
51}
52
53pub(super) fn handle_server_request(
54    core: &mut ProtocolCore,
55    request_id: RequestId,
56    request: ServerRequest,
57) -> Commands {
58    match request {
59        ServerRequest::Offer(payload) => {
60            handle_negotiation_request(core, request_id, NegotiationKind::Offer, payload)
61        }
62        ServerRequest::Renegotiate(payload) => {
63            handle_negotiation_request(core, request_id, NegotiationKind::Renegotiate, payload)
64        }
65    }
66}
67
68pub(super) fn handle_server_response(
69    core: &mut ProtocolCore,
70    response_to: &RequestId,
71    response: ServerResponse,
72) -> Commands {
73    match response {
74        ServerResponse::StartRecording(payload) => core.request_tracker.resolve_response(
75            response_to,
76            PendingRequestKind::StartRecording,
77            payload.ok,
78        ),
79        ServerResponse::StopRecording(payload) => core.request_tracker.resolve_response(
80            response_to,
81            PendingRequestKind::StopRecording,
82            payload.ok,
83        ),
84    }
85}
86
87fn handle_negotiation_request(
88    core: &mut ProtocolCore,
89    request_id: RequestId,
90    kind: NegotiationKind,
91    payload: SessionDescriptionPayload,
92) -> Commands {
93    match core.phase.accept_negotiation(&request_id, kind) {
94        Ok(()) => {}
95        Err(NegotiationRejection::Ignored) => return Vec::new(),
96        Err(NegotiationRejection::ProtocolError) => {
97            return close_for_protocol_error();
98        }
99    }
100    vec![Command::ApplyNegotiation {
101        request_id,
102        kind,
103        sdp: payload.sdp,
104        upload_slots: payload.upload_slots,
105    }]
106}
107
108fn begin_request(core: &mut ProtocolCore, request: ClientRequest) -> Commands {
109    if !core.phase.can_send_client_messages() {
110        return Vec::new();
111    }
112    let kind = match &request {
113        ClientRequest::StartRecording(_) => PendingRequestKind::StartRecording,
114        ClientRequest::StopRecording => PendingRequestKind::StopRecording,
115    };
116    let Some(request_start) = core.request_tracker.try_begin(kind) else {
117        return Vec::new();
118    };
119    let request_id = request_start.request_id;
120    let pending_request = PendingRequest {
121        request_id: request_id.clone(),
122        timeout_timer_id: request_start.timeout_timer_id.raw(),
123        timeout_ms: REQUEST_TIMEOUT_MS,
124    };
125    let Some(envelope) = ClientEnvelope::Request {
126        request_id,
127        request,
128    }
129    .into_envelope()
130    .ok() else {
131        return Vec::new();
132    };
133    let mut commands = vec![Command::BeginPendingRequest {
134        request: pending_request,
135    }];
136    commands.extend(core.outbound_batch.enqueue(envelope, FlushMode::Batched));
137    commands
138}