1use std::sync::Arc;
4
5use futures_util::Stream;
6use reqwest::Method;
7use serde::{Deserialize, Serialize};
8use serde_json::{Map, Value};
9
10use crate::client::{CallOptions, Client};
11use crate::common::{
12 CompositionConfig, LayoutPriority, ResourceLinks, TranscriptionConfig, WebhookDeliverySummary,
13};
14use crate::error::{Error, Result};
15use crate::pagination::{auto_page, paginate, ListParams, Page, PageFetcher};
16use crate::query::QueryBuilder;
17use crate::resources::egress::{
18 composition_to_config, to_egress_handle, Composition, EgressHandle, EgressType, StopTarget,
19 StopWire,
20};
21use crate::resources::escape;
22
23const PATH: &str = "/v2/hls";
24
25#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
27#[serde(rename_all = "lowercase")]
28pub enum HlsCaptureFormat {
29 Png,
31 Jpg,
33 Webp,
35}
36
37#[derive(Debug, Clone, Deserialize)]
39#[serde(rename_all = "camelCase")]
40pub struct HlsStream {
41 pub id: String,
43 pub composer_id: Option<String>,
45 pub room_id: Option<String>,
47 pub session_id: Option<String>,
49 pub mode: Option<String>,
51 pub start: Option<String>,
53 pub end: Option<String>,
55 pub quality: Option<String>,
57 pub orientation: Option<String>,
59 #[serde(default, deserialize_with = "crate::common::null_to_default")]
61 pub capture_log: Vec<Value>,
62 pub encryption: Option<Map<String, Value>>,
64 pub duration: Option<f64>,
66 pub webhook: Option<WebhookDeliverySummary>,
68 #[serde(default, deserialize_with = "crate::common::null_to_default")]
70 pub links: ResourceLinks,
71 pub downstream_url: Option<String>,
73 pub playback_hls_url: Option<String>,
75 pub livestream_url: Option<String>,
77 #[serde(flatten)]
79 pub extra: Map<String, Value>,
80}
81
82#[derive(Debug, Clone, Default)]
84pub struct HlsStartParams {
85 pub composition: Option<Composition>,
87 pub transcription: Option<TranscriptionConfig>,
89 pub webhook_url: Option<String>,
91 pub recording: Option<Value>,
93 pub template_url: Option<String>,
98 pub resource_id: Option<String>,
100}
101
102#[derive(Debug, Serialize)]
103#[serde(rename_all = "camelCase")]
104struct HlsStartWire<'a> {
105 room_id: &'a str,
106 #[serde(skip_serializing_if = "Option::is_none")]
107 config: Option<CompositionConfig>,
108 #[serde(skip_serializing_if = "Option::is_none")]
109 template_url: Option<String>,
110 #[serde(skip_serializing_if = "Option::is_none")]
111 transcription: Option<&'a TranscriptionConfig>,
112 #[serde(skip_serializing_if = "Option::is_none")]
113 webhook_url: Option<&'a str>,
114 #[serde(skip_serializing_if = "Option::is_none")]
115 recording: Option<&'a Value>,
116 #[serde(skip_serializing_if = "Option::is_none")]
117 resource_id: Option<&'a str>,
118}
119
120#[derive(Debug, Clone, Default, Serialize)]
122#[serde(rename_all = "camelCase")]
123pub struct HlsCaptureParams {
124 pub room_id: String,
126 #[serde(skip_serializing_if = "Option::is_none")]
128 pub format: Option<HlsCaptureFormat>,
129 #[serde(skip_serializing_if = "Option::is_none")]
131 pub width: Option<u32>,
132 #[serde(skip_serializing_if = "Option::is_none")]
134 pub height: Option<u32>,
135 #[serde(skip_serializing_if = "Option::is_none")]
137 pub time: Option<u32>,
138}
139
140#[derive(Debug, Clone, Deserialize)]
142#[serde(rename_all = "camelCase")]
143pub struct HlsCaptureResult {
144 pub message: String,
146 pub room_id: String,
148 pub file_path: Option<String>,
150 pub file_size: Option<i64>,
152 pub file_name: Option<String>,
154 pub meta: Option<Value>,
156 #[serde(flatten)]
158 pub extra: Map<String, Value>,
159}
160
161#[derive(Debug, Clone, Default)]
163pub struct HlsListParams {
164 pub page: Option<u32>,
166 pub per_page: Option<u32>,
168 pub cursor: Option<String>,
170 pub room_id: Option<String>,
172 pub session_id: Option<String>,
174 pub user_id: Option<String>,
176 pub query: Option<String>,
178}
179
180impl HlsListParams {
181 fn pagination(&self) -> ListParams {
182 ListParams {
183 page: self.page,
184 per_page: self.per_page,
185 cursor: self.cursor.clone(),
186 }
187 }
188}
189
190#[derive(Debug, Clone, Copy)]
192pub struct HlsResource<'a> {
193 client: &'a Client,
194}
195
196impl<'a> HlsResource<'a> {
197 pub(crate) fn new(client: &'a Client) -> Self {
198 Self { client }
199 }
200
201 pub async fn start(&self, room_id: &str, params: HlsStartParams) -> Result<EgressHandle> {
203 let mapped = composition_to_config(
204 params.composition.as_ref(),
205 Some(LayoutPriority::Speaker),
207 );
208 let template_url = mapped.template_url.or(params.template_url);
210
211 let body = HlsStartWire {
212 room_id,
213 config: mapped.config,
214 template_url,
215 transcription: params.transcription.as_ref(),
216 webhook_url: params.webhook_url.as_deref(),
217 recording: params.recording.as_ref(),
218 resource_id: params.resource_id.as_deref(),
219 };
220 let path = format!("{PATH}/start");
221 let raw = self
222 .client
223 .maybe_json(Method::POST, &path, CallOptions::json(&body)?)
224 .await?;
225 Ok(to_egress_handle(EgressType::Hls, room_id, raw))
226 }
227
228 pub async fn stop(&self, target: impl Into<StopTarget>) -> Result<String> {
230 let target = target.into();
231 let path = format!("{PATH}/end");
232 self.client
233 .message(
234 Method::POST,
235 &path,
236 CallOptions::json(StopWire::from(&target))?,
237 )
238 .await
239 }
240
241 pub async fn get(&self, room_id: &str) -> Result<HlsStream> {
245 let path = format!("{PATH}/{}/active", escape(room_id));
246 let raw: Value = self
247 .client
248 .json(Method::GET, &path, CallOptions::new())
249 .await?;
250
251 let data = raw.get("data").filter(|data| !data.is_null());
252 let Some(data) = data else {
253 return Err(Error::not_found(format!(
254 "no active HLS stream for room {room_id}"
255 )));
256 };
257
258 let mut stream: HlsStream =
259 serde_json::from_value(data.clone()).map_err(|e| Error::decode(path, e))?;
260 if stream.room_id.is_none() {
262 stream.room_id = stream
263 .extra
264 .get("meetingId")
265 .and_then(Value::as_str)
266 .map(str::to_string);
267 }
268 Ok(stream)
269 }
270
271 pub async fn get_by_id(&self, id: &str) -> Result<HlsStream> {
273 let path = format!("{PATH}/{}", escape(id));
274 self.client
275 .json(Method::GET, &path, CallOptions::new())
276 .await
277 }
278
279 pub async fn capture(&self, params: HlsCaptureParams) -> Result<HlsCaptureResult> {
281 let path = format!("{PATH}/capture");
282 self.client
283 .json(Method::POST, &path, CallOptions::json(¶ms)?)
284 .await
285 }
286
287 pub async fn list(&self, params: HlsListParams) -> Result<Page<HlsStream>> {
289 paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
290 }
291
292 pub fn list_stream(
294 &self,
295 params: HlsListParams,
296 ) -> impl Stream<Item = Result<HlsStream>> + Send {
297 auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
298 }
299
300 pub async fn delete(&self, id: &str) -> Result<String> {
302 let path = format!("{PATH}/{}", escape(id));
303 self.client
304 .message(Method::DELETE, &path, CallOptions::new())
305 .await
306 }
307
308 fn fetcher(&self, params: &HlsListParams) -> PageFetcher {
309 let client = self.client.clone();
310 let params = params.clone();
311 Arc::new(move |page, per_page| {
312 let client = client.clone();
313 let params = params.clone();
314 Box::pin(async move {
315 let query = QueryBuilder::new()
316 .opt("page", page)
317 .opt("perPage", per_page)
318 .opt_str("roomId", params.room_id.as_deref())
319 .opt_str("sessionId", params.session_id.as_deref())
320 .opt_str("userId", params.user_id.as_deref())
321 .opt_str("query", params.query.as_deref())
322 .into_pairs();
323 client
324 .json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
325 .await
326 })
327 })
328 }
329}