use std::sync::Arc;
use futures_util::Stream;
use reqwest::Method;
use serde::{Deserialize, Serialize};
use serde_json::{Map, Value};
use crate::client::{CallOptions, Client};
use crate::common::{
CompositionConfig, OnFailureConfig, ParticipantFileFormat, RecordingFile, RecordingFileFormat,
ResourceLinks, TrackFileFormat, TrackKind, TranscriptionConfig, WebhookDeliverySummary,
};
use crate::error::{Error, Result};
use crate::pagination::{auto_page, paginate, ListParams, Page, PageFetcher};
use crate::query::QueryBuilder;
use crate::resources::egress::{
composition_to_config, to_egress_handle, Composition, EgressHandle, EgressType, StopTarget,
StopWire,
};
use crate::resources::escape;
const PATH: &str = "/v2/recordings";
const PARTICIPANT_PATH: &str = "/v2/recordings/participant";
const TRACK_PATH: &str = "/v2/recordings/participant/track";
const COMPOSITE_PATH: &str = "/v2/recordings/composite";
const MERGE_PATH: &str = "/v2/recordings/participant/merge";
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Recording {
pub id: String,
pub room_id: Option<String>,
pub file: Option<RecordingFile>,
pub webhook: Option<WebhookDeliverySummary>,
pub start: Option<String>,
pub end: Option<String>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub links: ResourceLinks,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct IndividualRecording {
pub id: String,
pub room_id: Option<String>,
pub participant_name: Option<String>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub files: Vec<RecordingFile>,
pub webhook: Option<WebhookDeliverySummary>,
pub start: Option<String>,
pub end: Option<String>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct CompositeRecording {
pub id: String,
pub meeting_id: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub participants: Vec<Value>,
pub file_format: Option<String>,
pub start: Option<String>,
pub end: Option<String>,
pub file_id: Option<String>,
pub file: Option<RecordingFile>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct MergeChannel {
pub recording_id: Option<String>,
pub participant_id: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct MergeRecording {
pub id: String,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub channel1: Vec<MergeChannel>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub channel2: Vec<MergeChannel>,
pub start: Option<String>,
pub end: Option<String>,
pub status: Option<String>,
pub meeting_id: Option<String>,
pub session_id: Option<String>,
pub file: Option<RecordingFile>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct MergeRecordingResult {
pub message: String,
pub recording: MergeRecording,
}
#[derive(Debug, Clone, Default)]
pub struct RecordingStartParams {
pub composition: Option<Composition>,
pub transcription: Option<TranscriptionConfig>,
pub on_failure: Option<OnFailureConfig>,
pub file_format: Option<RecordingFileFormat>,
pub webhook_url: Option<String>,
pub dir_path: Option<String>,
pub resource_id: Option<String>,
pub pre_signed_url: Option<String>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct RecordingStartWire<'a> {
room_id: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
config: Option<CompositionConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
template_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
aws_dir_path: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
transcription: Option<&'a TranscriptionConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
on_failure: Option<&'a OnFailureConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
file_format: Option<RecordingFileFormat>,
#[serde(skip_serializing_if = "Option::is_none")]
webhook_url: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
resource_id: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
pre_signed_url: Option<&'a str>,
}
#[derive(Debug, Clone, Default)]
pub struct RecordingGetParams {
pub with_transcription: Option<bool>,
}
#[derive(Debug, Clone, Default)]
pub struct ParticipantRecordingStartParams {
pub participant_id: String,
pub file_format: Option<ParticipantFileFormat>,
pub webhook_url: Option<String>,
pub dir_path: Option<String>,
pub metadata: Option<Value>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct ParticipantRecordingStartWire<'a> {
room_id: &'a str,
participant_id: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
file_format: Option<ParticipantFileFormat>,
#[serde(skip_serializing_if = "Option::is_none")]
webhook_url: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
bucket_dir_path: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<&'a Value>,
}
#[derive(Debug, Clone)]
pub struct TrackRecordingStartParams {
pub participant_id: String,
pub kind: TrackKind,
pub file_format: Option<TrackFileFormat>,
pub webhook_url: Option<String>,
pub dir_path: Option<String>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct TrackRecordingStartWire<'a> {
room_id: &'a str,
participant_id: &'a str,
kind: TrackKind,
#[serde(skip_serializing_if = "Option::is_none")]
file_format: Option<TrackFileFormat>,
#[serde(skip_serializing_if = "Option::is_none")]
webhook_url: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
bucket_dir_path: Option<&'a str>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CompositeParticipantSelector {
pub participant_id: String,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub kind: Vec<TrackKind>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct CompositeWatermark {
#[serde(rename = "type")]
pub kind: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub image_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub image_base64: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub text: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub timezone: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub topic: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct CompositeRecordingStartParams {
pub participants: Vec<CompositeParticipantSelector>,
pub watermarks: Vec<CompositeWatermark>,
pub webhook_url: Option<String>,
pub dir_path: Option<String>,
pub pre_signed_url: Option<String>,
pub file_format: Option<String>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct CompositeRecordingStartWire<'a> {
room_id: &'a str,
#[serde(skip_serializing_if = "<[_]>::is_empty")]
participants: &'a [CompositeParticipantSelector],
#[serde(skip_serializing_if = "<[_]>::is_empty")]
watermarks: &'a [CompositeWatermark],
#[serde(skip_serializing_if = "Option::is_none")]
webhook_url: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
bucket_dir_path: Option<&'a str>,
#[serde(rename = "preSignedURL", skip_serializing_if = "Option::is_none")]
pre_signed_url: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
file_format: Option<&'a str>,
}
#[derive(Debug, Clone, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct MergeChannelEntry {
pub participant_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub recording_id: Option<String>,
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
pub kind: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct MergeRecordingCreateParams {
pub session_id: String,
pub channel1: Vec<MergeChannelEntry>,
pub channel2: Vec<MergeChannelEntry>,
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
pub kind: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub webhook_url: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub dir_path: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct RecordingListParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
pub user_id: Option<String>,
pub query: Option<String>,
pub with_transcription: Option<bool>,
}
#[derive(Debug, Clone, Default)]
pub struct IndividualRecordingListParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
pub participant_id: Option<String>,
pub user_id: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct CompositeRecordingListParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct MergeRecordingListParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub status: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
pub id: Option<String>,
}
macro_rules! pagination {
($name:ident) => {
impl $name {
fn pagination(&self) -> ListParams {
ListParams {
page: self.page,
per_page: self.per_page,
cursor: self.cursor.clone(),
}
}
}
};
}
pagination!(RecordingListParams);
pagination!(IndividualRecordingListParams);
pagination!(CompositeRecordingListParams);
pagination!(MergeRecordingListParams);
#[derive(Debug, Clone, Copy)]
pub struct RecordingParticipantResource<'a> {
client: &'a Client,
}
impl<'a> RecordingParticipantResource<'a> {
pub async fn start(
&self,
room_id: &str,
params: ParticipantRecordingStartParams,
) -> Result<String> {
let body = ParticipantRecordingStartWire {
room_id,
participant_id: ¶ms.participant_id,
file_format: params.file_format,
webhook_url: params.webhook_url.as_deref(),
bucket_dir_path: params.dir_path.as_deref(),
metadata: params.metadata.as_ref(),
};
let path = format!("{PARTICIPANT_PATH}/start");
self.client
.message(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn stop(&self, room_id: &str, participant_id: &str) -> Result<String> {
let body = serde_json::json!({"roomId": room_id, "participantId": participant_id});
let path = format!("{PARTICIPANT_PATH}/stop");
self.client
.message(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn list(
&self,
params: IndividualRecordingListParams,
) -> Result<Page<IndividualRecording>> {
let fetcher = individual_fetcher(self.client, PARTICIPANT_PATH, ¶ms);
paginate(fetcher, ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: IndividualRecordingListParams,
) -> impl Stream<Item = Result<IndividualRecording>> + Send {
let fetcher = individual_fetcher(self.client, PARTICIPANT_PATH, ¶ms);
auto_page(fetcher, params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<IndividualRecording> {
let path = format!("{PARTICIPANT_PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
pub async fn delete(&self, id: &str) -> Result<String> {
let path = format!("{PARTICIPANT_PATH}/{}", escape(id));
self.client
.message(Method::DELETE, &path, CallOptions::new())
.await
}
}
#[derive(Debug, Clone, Copy)]
pub struct RecordingTrackResource<'a> {
client: &'a Client,
}
impl<'a> RecordingTrackResource<'a> {
pub async fn start(&self, room_id: &str, params: TrackRecordingStartParams) -> Result<String> {
let body = TrackRecordingStartWire {
room_id,
participant_id: ¶ms.participant_id,
kind: params.kind,
file_format: params.file_format,
webhook_url: params.webhook_url.as_deref(),
bucket_dir_path: params.dir_path.as_deref(),
};
let path = format!("{TRACK_PATH}/start");
self.client
.message(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn stop(
&self,
room_id: &str,
participant_id: &str,
kind: TrackKind,
) -> Result<String> {
let body = serde_json::json!({
"roomId": room_id, "participantId": participant_id, "kind": kind,
});
let path = format!("{TRACK_PATH}/stop");
self.client
.message(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn list(
&self,
params: IndividualRecordingListParams,
) -> Result<Page<IndividualRecording>> {
let fetcher = individual_fetcher(self.client, TRACK_PATH, ¶ms);
paginate(fetcher, ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: IndividualRecordingListParams,
) -> impl Stream<Item = Result<IndividualRecording>> + Send {
let fetcher = individual_fetcher(self.client, TRACK_PATH, ¶ms);
auto_page(fetcher, params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<IndividualRecording> {
let path = format!("{TRACK_PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
pub async fn delete(&self, id: &str) -> Result<String> {
let path = format!("{TRACK_PATH}/{}", escape(id));
self.client
.message(Method::DELETE, &path, CallOptions::new())
.await
}
}
#[derive(Debug, Clone, Copy)]
pub struct RecordingCompositeResource<'a> {
client: &'a Client,
}
impl<'a> RecordingCompositeResource<'a> {
pub async fn start(
&self,
room_id: &str,
params: CompositeRecordingStartParams,
) -> Result<EgressHandle> {
let body = CompositeRecordingStartWire {
room_id,
participants: ¶ms.participants,
watermarks: ¶ms.watermarks,
webhook_url: params.webhook_url.as_deref(),
bucket_dir_path: params.dir_path.as_deref(),
pre_signed_url: params.pre_signed_url.as_deref(),
file_format: params.file_format.as_deref(),
};
let path = format!("{COMPOSITE_PATH}/start");
let raw = self
.client
.maybe_json(Method::POST, &path, CallOptions::json(&body)?)
.await?;
Ok(to_egress_handle(EgressType::Composite, room_id, raw))
}
pub async fn stop(&self, handle: &EgressHandle) -> Result<String> {
let id = handle
.id
.as_deref()
.filter(|id| !id.is_empty())
.ok_or_else(|| {
Error::validation(
"recordings.composite().stop() requires a recordingId, from the start handle",
)
})?;
let body = serde_json::json!({"roomId": handle.room_id, "recordingId": id});
let path = format!("{COMPOSITE_PATH}/stop");
self.client
.message(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn list(
&self,
params: CompositeRecordingListParams,
) -> Result<Page<CompositeRecording>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: CompositeRecordingListParams,
) -> impl Stream<Item = Result<CompositeRecording>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<CompositeRecording> {
let path = format!("{COMPOSITE_PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
pub async fn delete(&self, id: &str) -> Result<String> {
let path = format!("{COMPOSITE_PATH}/{}", escape(id));
self.client
.message(Method::DELETE, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &CompositeRecordingListParams) -> 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())
.into_pairs();
client
.json::<Value>(Method::GET, COMPOSITE_PATH, CallOptions::new().query(query))
.await
})
})
}
}
#[derive(Debug, Clone, Copy)]
pub struct MergeRecordingResource<'a> {
client: &'a Client,
}
impl<'a> MergeRecordingResource<'a> {
pub async fn create(&self, params: MergeRecordingCreateParams) -> Result<MergeRecordingResult> {
self.client
.json(Method::POST, MERGE_PATH, CallOptions::json(¶ms)?)
.await
}
pub async fn list(&self, params: MergeRecordingListParams) -> Result<Page<MergeRecording>> {
paginate(
self.fetcher(¶ms),
¶ms.pagination(),
"recordings",
None,
)
.await
}
pub fn list_stream(
&self,
params: MergeRecordingListParams,
) -> impl Stream<Item = Result<MergeRecording>> + Send {
auto_page(
self.fetcher(¶ms),
params.pagination(),
"recordings",
None,
)
}
pub async fn get(&self, muxer_id: &str) -> Result<MergeRecording> {
let path = format!("{MERGE_PATH}/{}", escape(muxer_id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &MergeRecordingListParams) -> 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("status", params.status.as_deref())
.opt_str("roomId", params.room_id.as_deref())
.opt_str("sessionId", params.session_id.as_deref())
.opt_str("id", params.id.as_deref())
.into_pairs();
client
.json::<Value>(Method::GET, MERGE_PATH, CallOptions::new().query(query))
.await
})
})
}
}
fn individual_fetcher(
client: &Client,
path: &'static str,
params: &IndividualRecordingListParams,
) -> PageFetcher {
let client = 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("participantId", params.participant_id.as_deref())
.opt_str("userId", params.user_id.as_deref())
.into_pairs();
client
.json::<Value>(Method::GET, path, CallOptions::new().query(query))
.await
})
})
}
#[derive(Debug, Clone, Copy)]
pub struct RecordingsResource<'a> {
client: &'a Client,
}
impl<'a> RecordingsResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub fn participant(&self) -> RecordingParticipantResource<'a> {
RecordingParticipantResource {
client: self.client,
}
}
pub fn track(&self) -> RecordingTrackResource<'a> {
RecordingTrackResource {
client: self.client,
}
}
pub fn composite(&self) -> RecordingCompositeResource<'a> {
RecordingCompositeResource {
client: self.client,
}
}
pub fn merge(&self) -> MergeRecordingResource<'a> {
MergeRecordingResource {
client: self.client,
}
}
pub async fn start(&self, room_id: &str, params: RecordingStartParams) -> Result<EgressHandle> {
let mapped = composition_to_config(params.composition.as_ref(), None);
let body = RecordingStartWire {
room_id,
config: mapped.config,
template_url: mapped.template_url,
aws_dir_path: params.dir_path.as_deref(),
transcription: params.transcription.as_ref(),
on_failure: params.on_failure.as_ref(),
file_format: params.file_format,
webhook_url: params.webhook_url.as_deref(),
resource_id: params.resource_id.as_deref(),
pre_signed_url: params.pre_signed_url.as_deref(),
};
let path = format!("{PATH}/start");
let raw = self
.client
.maybe_json(Method::POST, &path, CallOptions::json(&body)?)
.await?;
Ok(to_egress_handle(EgressType::Recording, room_id, raw))
}
pub async fn stop(&self, target: impl Into<StopTarget>) -> Result<String> {
let target = target.into();
let path = format!("{PATH}/end");
self.client
.message(
Method::POST,
&path,
CallOptions::json(StopWire::from(&target))?,
)
.await
}
pub async fn list(&self, params: RecordingListParams) -> Result<Page<Recording>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: RecordingListParams,
) -> impl Stream<Item = Result<Recording>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, id: &str, params: RecordingGetParams) -> Result<Recording> {
let query = QueryBuilder::new()
.opt("withTranscription", params.with_transcription)
.into_pairs();
let path = format!("{PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new().query(query))
.await
}
pub async fn delete(&self, id: &str) -> Result<String> {
let path = format!("{PATH}/{}", escape(id));
self.client
.message(Method::DELETE, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &RecordingListParams) -> 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("userId", params.user_id.as_deref())
.opt_str("query", params.query.as_deref())
.opt("withTranscription", params.with_transcription)
.into_pairs();
client
.json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
.await
})
})
}
}