use std::sync::Arc;
use futures_util::Stream;
use reqwest::Method;
use serde::{Deserialize, Serialize, Serializer};
use serde_json::{Map, Value};
use crate::client::{CallOptions, Client};
use crate::common::{string_enum, WebhookDeliverySummary};
use crate::error::Result;
use crate::pagination::{auto_page, paginate, ListParams, Page, PageFetcher};
use crate::query::QueryBuilder;
use crate::resources::escape;
const PATH: &str = "/v2/transcodings";
string_enum! {
TranscodingStatus {
PENDING => "pending",
PROCESSING => "processing",
COMPLETED => "completed",
FAILED => "failed",
CANCELLED => "cancelled",
}
}
string_enum! {
TranscodingTask {
COMPOSITE_MERGE => "composite-merge",
HLS_TO_MP4 => "hls-to-mp4",
MEETING_RECORDING_MERGE => "meeting-recording-merge",
}
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct TranscodingWatermark {
#[serde(rename = "type")]
pub kind: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub image: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub start_time: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub lat: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub long: Option<f64>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct TranscodingStorage {
#[serde(rename = "type")]
pub kind: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub bucket: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub container: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub dir_path: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct MergeTranscodingParams {
pub recording_ids: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub task: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub webhook_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub watermark: Option<TranscodingWatermark>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct HlsToMp4Params {
#[serde(skip_serializing_if = "Option::is_none")]
pub room_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub session_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub hls_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub webhook_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub storage: Option<TranscodingStorage>,
}
#[derive(Debug, Clone)]
pub struct MeetingRecordingRef {
pub id: String,
pub presigned_url: Option<String>,
}
impl MeetingRecordingRef {
pub fn new(id: impl Into<String>) -> Self {
Self {
id: id.into(),
presigned_url: None,
}
}
}
impl From<&str> for MeetingRecordingRef {
fn from(id: &str) -> Self {
Self::new(id)
}
}
impl From<String> for MeetingRecordingRef {
fn from(id: String) -> Self {
Self::new(id)
}
}
impl Serialize for MeetingRecordingRef {
fn serialize<S: Serializer>(&self, serializer: S) -> std::result::Result<S::Ok, S::Error> {
match &self.presigned_url {
None => serializer.serialize_str(&self.id),
Some(presigned_url) => {
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct Full<'a> {
id: &'a str,
presigned_url: &'a str,
}
Full {
id: &self.id,
presigned_url,
}
.serialize(serializer)
}
}
}
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct MeetingRecordingMergeParams {
pub recording_ids: Vec<MeetingRecordingRef>,
#[serde(skip_serializing_if = "Option::is_none")]
pub webhook_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub presigned_output_url: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct TranscodingFile {
pub id: Option<String>,
#[serde(rename = "type")]
pub kind: Option<String>,
pub size: Option<i64>,
pub meta: Option<Value>,
pub file_path: Option<String>,
pub file_url: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Transcoding {
pub id: String,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub recording_ids: Vec<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
pub hls_id: Option<String>,
pub status: Option<TranscodingStatus>,
pub task: Option<TranscodingTask>,
pub started_at: Option<String>,
pub stopped_at: Option<String>,
pub file: Option<TranscodingFile>,
pub webhook: Option<WebhookDeliverySummary>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Default)]
pub struct ListTranscodingsParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
pub hls_id: Option<String>,
pub status: Option<TranscodingStatus>,
}
impl ListTranscodingsParams {
fn pagination(&self) -> ListParams {
ListParams {
page: self.page,
per_page: self.per_page,
cursor: self.cursor.clone(),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct TranscodingsResource<'a> {
client: &'a Client,
}
impl<'a> TranscodingsResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub async fn merge(&self, params: MergeTranscodingParams) -> Result<Transcoding> {
let path = format!("{PATH}/merge");
self.client
.json(Method::POST, &path, CallOptions::json(¶ms)?)
.await
}
pub async fn hls_to_mp4(&self, params: HlsToMp4Params) -> Result<Transcoding> {
let path = format!("{PATH}/hls-to-mp4");
self.client
.json(Method::POST, &path, CallOptions::json(¶ms)?)
.await
}
pub async fn meeting_recording_merge(
&self,
params: MeetingRecordingMergeParams,
) -> Result<Transcoding> {
let path = format!("{PATH}/meeting-recording-merge");
self.client
.json(Method::POST, &path, CallOptions::json(¶ms)?)
.await
}
pub async fn list(&self, params: ListTranscodingsParams) -> Result<Page<Transcoding>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: ListTranscodingsParams,
) -> impl Stream<Item = Result<Transcoding>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<Transcoding> {
let path = format!("{PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
pub async fn cancel(&self, id: &str) -> Result<Transcoding> {
let path = format!("{PATH}/{}/cancel", escape(id));
self.client
.json(Method::POST, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &ListTranscodingsParams) -> PageFetcher {
let client = self.client.clone();
let params = params.clone();
Arc::new(move |page, per_page| {
let client = client.clone();
let params = params.clone();
Box::pin(async move {
let query = QueryBuilder::new()
.opt("page", page)
.opt("perPage", per_page)
.opt_str("roomId", params.room_id.as_deref())
.opt_str("sessionId", params.session_id.as_deref())
.opt_str("hlsId", params.hls_id.as_deref())
.opt_str(
"status",
params.status.as_ref().map(TranscodingStatus::as_str),
)
.into_pairs();
client
.json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
.await
})
})
}
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn a_bare_recording_ref_serializes_as_a_string() {
let params = MeetingRecordingMergeParams {
recording_ids: vec!["rec-1".into(), "rec-2".into()],
..Default::default()
};
assert_eq!(
serde_json::to_value(¶ms).unwrap(),
json!({"recordingIds": ["rec-1", "rec-2"]})
);
}
#[test]
fn a_recording_ref_with_a_presigned_url_serializes_as_an_object() {
let params = MeetingRecordingMergeParams {
recording_ids: vec![
MeetingRecordingRef::new("rec-1"),
MeetingRecordingRef {
id: "rec-2".into(),
presigned_url: Some("https://s3/get".into()),
},
],
presigned_output_url: Some("https://s3/put".into()),
..Default::default()
};
assert_eq!(
serde_json::to_value(¶ms).unwrap(),
json!({
"recordingIds": ["rec-1", {"id": "rec-2", "presignedUrl": "https://s3/get"}],
"presignedOutputUrl": "https://s3/put",
})
);
}
}