o_sfu_protocol/core/
request_flow.rs1use 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 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}