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, LayoutPriority, ResourceLinks, 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/livestreams";
#[derive(Debug, Clone, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct RtmpStream {
pub url: String,
pub stream_key: String,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct Livestream {
pub id: String,
pub composer_id: Option<String>,
pub room_id: Option<String>,
pub session_id: Option<String>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub outputs: Vec<RtmpStream>,
pub start: Option<String>,
pub end: Option<String>,
pub region: Option<String>,
pub webhook: Option<WebhookDeliverySummary>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub links: ResourceLinks,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Default)]
pub struct RtmpStartParams {
pub streams: Vec<RtmpStream>,
pub composition: Option<Composition>,
pub transcription: Option<TranscriptionConfig>,
pub resource_id: Option<String>,
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct RtmpStartWire<'a> {
room_id: &'a str,
outputs: &'a [RtmpStream],
#[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")]
transcription: Option<&'a TranscriptionConfig>,
#[serde(skip_serializing_if = "Option::is_none")]
resource_id: Option<&'a str>,
}
#[derive(Debug, Clone, Default)]
pub struct RtmpListParams {
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>,
}
impl RtmpListParams {
fn pagination(&self) -> ListParams {
ListParams {
page: self.page,
per_page: self.per_page,
cursor: self.cursor.clone(),
}
}
}
#[derive(Debug, Clone, Copy)]
pub struct RtmpResource<'a> {
client: &'a Client,
}
impl<'a> RtmpResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub async fn start(&self, room_id: &str, params: RtmpStartParams) -> Result<EgressHandle> {
if params.streams.is_empty() {
return Err(Error::validation(
"rtmp.start() requires at least one destination in `streams`",
));
}
let mapped = composition_to_config(
params.composition.as_ref(),
Some(LayoutPriority::Speaker),
);
let body = RtmpStartWire {
room_id,
outputs: ¶ms.streams,
config: mapped.config,
template_url: mapped.template_url,
transcription: params.transcription.as_ref(),
resource_id: params.resource_id.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::Livestream, 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: RtmpListParams) -> Result<Page<Livestream>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: RtmpListParams,
) -> impl Stream<Item = Result<Livestream>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<Livestream> {
let path = format!("{PATH}/{}", escape(id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
fn fetcher(&self, params: &RtmpListParams) -> 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())
.into_pairs();
client
.json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
.await
})
})
}
}