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::string_enum;
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/resource";
string_enum! {
ResourceStatus {
PENDING => "pending",
IDLE => "idle",
COMPOSING => "composing",
RELEASED => "released",
}
}
string_enum! {
ComposerType {
RECORDING => "recording",
HLS => "hls",
RTMP => "rtmp",
}
}
string_enum! {
ResourceMode {
VIDEO_AND_AUDIO => "video-and-audio",
AUDIO => "audio",
}
}
string_enum! {
ResourceQuality {
LOW => "low",
MED => "med",
HIGH => "high",
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct ResourceUnit {
pub id: String,
pub status: Option<ResourceStatus>,
#[serde(rename = "type")]
pub kind: Option<String>,
pub mode: Option<ResourceMode>,
pub quality: Option<ResourceQuality>,
#[serde(default, deserialize_with = "crate::common::null_to_default")]
pub composer_ids: Vec<String>,
pub webhook_url: Option<String>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Default)]
pub struct ListResourcesParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub status: Option<ResourceStatus>,
}
impl ListResourcesParams {
fn pagination(&self) -> ListParams {
ListParams {
page: self.page,
per_page: self.per_page,
cursor: self.cursor.clone(),
}
}
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AcquireResourceParams {
#[serde(rename = "type", skip_serializing_if = "Option::is_none")]
pub kind: Option<ComposerType>,
pub webhook_url: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub units: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub mode: Option<ResourceMode>,
#[serde(skip_serializing_if = "Option::is_none")]
pub quality: Option<ResourceQuality>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct ReleaseResourceResult {
pub id: String,
pub success: bool,
pub msg: Option<String>,
}
#[derive(Debug, Clone, Copy)]
pub struct ResourcePoolResource<'a> {
client: &'a Client,
}
impl<'a> ResourcePoolResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub async fn list(&self, params: ListResourcesParams) -> Result<Page<ResourceUnit>> {
paginate(self.fetcher(¶ms), ¶ms.pagination(), "data", None).await
}
pub fn list_stream(
&self,
params: ListResourcesParams,
) -> impl Stream<Item = Result<ResourceUnit>> + Send {
auto_page(self.fetcher(¶ms), params.pagination(), "data", None)
}
pub async fn get(&self, id: &str) -> Result<ResourceUnit> {
let path = format!("{PATH}/{}", escape(id));
self.client
.data(Method::GET, &path, CallOptions::new())
.await
}
pub async fn acquire(&self, params: AcquireResourceParams) -> Result<Vec<ResourceUnit>> {
let path = format!("{PATH}/acquire");
self.client
.data(Method::POST, &path, CallOptions::json(¶ms)?)
.await
}
pub async fn release(&self, ids: &[String]) -> Result<Vec<ReleaseResourceResult>> {
let body = serde_json::json!({ "ids": ids });
let path = format!("{PATH}/release");
self.client
.data(Method::POST, &path, CallOptions::json(&body)?)
.await
}
fn fetcher(&self, params: &ListResourcesParams) -> PageFetcher {
let client = self.client.clone();
let status = params.status.clone();
Arc::new(move |page, per_page| {
let client = client.clone();
let status = status.clone();
Box::pin(async move {
let query = QueryBuilder::new()
.opt("page", page)
.opt("perPage", per_page)
.opt_str("status", status.as_ref().map(ResourceStatus::as_str))
.into_pairs();
client
.json::<Value>(Method::GET, PATH, CallOptions::new().query(query))
.await
})
})
}
}