use std::collections::HashMap;
use std::sync::Arc;
use futures_util::Stream;
use reqwest::Method;
use serde::{Deserialize, Serialize};
use serde_json::{json, Map, Value};
use crate::client::{CallOptions, Client};
use crate::common::{string_enum, MessageResponse};
use crate::error::{Error, Result};
use crate::pagination::{auto_page, paginate, ListParams, Page, PageFetcher};
use crate::query::QueryBuilder;
use crate::resources::rooms::RoomWebhook;
use crate::resources::{escape, random_uuid, slugify_name};
const DEPLOYMENTS_PATH: &str = "/ai/v1/cloud/deployments";
const VERSIONS_PATH: &str = "/ai/v1/cloud/versions";
const SECRETS_PATH: &str = "/ai/v1/cloud/secrets";
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct AgentRoomOptions {
#[serde(skip_serializing_if = "Option::is_none")]
pub auto_end_session: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub session_timeout_seconds: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub playground: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub vision: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub join_meeting: Option<bool>,
}
#[derive(Debug, Clone, Default)]
pub struct DispatchAgentParams {
pub agent_id: Option<String>,
pub meeting_id: Option<String>,
pub room_id: Option<String>,
pub room_options: Option<AgentRoomOptions>,
pub sip_options: Option<Map<String, Value>>,
pub metadata: Option<Map<String, Value>>,
pub version_id: Option<String>,
pub version_tag: Option<String>,
pub wait: Option<Value>,
}
#[derive(Debug, Clone, Default)]
pub struct GeneralDispatchAgentParams {
pub agent_id: Option<String>,
pub meeting_id: Option<String>,
pub room_id: Option<String>,
pub room_options: Option<AgentRoomOptions>,
pub sip_options: Option<Map<String, Value>>,
pub metadata: Option<Map<String, Value>>,
pub version_id: Option<String>,
}
impl From<GeneralDispatchAgentParams> for DispatchAgentParams {
fn from(params: GeneralDispatchAgentParams) -> Self {
Self {
agent_id: params.agent_id,
meeting_id: params.meeting_id,
room_id: params.room_id,
room_options: params.room_options,
sip_options: params.sip_options,
metadata: params.metadata,
version_id: params.version_id,
version_tag: None,
wait: None,
}
}
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct DispatchWire {
#[serde(skip_serializing_if = "Option::is_none")]
agent_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
room_options: Option<AgentRoomOptions>,
#[serde(skip_serializing_if = "Option::is_none")]
sip_options: Option<Map<String, Value>>,
#[serde(skip_serializing_if = "Option::is_none")]
metadata: Option<Map<String, Value>>,
#[serde(skip_serializing_if = "Option::is_none")]
version_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
version_tag: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
wait: Option<Value>,
meeting_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
room_id: Option<String>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct DispatchResult {
pub success: Option<bool>,
pub job_id: Option<String>,
pub status: Option<String>,
pub worker_id: Option<String>,
pub room_id: Option<String>,
pub agent_id: Option<String>,
pub cloud: Option<Value>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AgentInitConfig {
pub registry_url: String,
}
string_enum! {
AgentComputeProfile {
CPU_SMALL => "cpu-small",
CPU_MEDIUM => "cpu-medium",
CPU_LARGE => "cpu-large",
}
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AgentDeployment {
pub id: String,
pub agent_id: Option<String>,
pub name: Option<String>,
pub template: Option<String>,
pub status: Option<String>,
pub env_secret: Option<String>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
#[serde(rename_all = "camelCase")]
pub struct AgentDeploymentVersion {
pub id: String,
pub deployment_id: Option<String>,
pub agent_id: Option<String>,
pub name: Option<String>,
pub image: Option<Value>,
pub region: Option<String>,
pub min_replica: Option<u32>,
pub max_replica: Option<u32>,
pub status: Option<String>,
pub active: Option<bool>,
pub version_tag: Option<String>,
#[serde(flatten)]
pub extra: Map<String, Value>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct AgentLogEntry {
pub log: String,
pub timestamp: Option<String>,
pub metadata: Option<Map<String, Value>>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct DeployImage {
pub image_url: String,
#[serde(rename = "imageCR", skip_serializing_if = "Option::is_none")]
pub image_cr: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub image_type: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub image_pull_secret: Option<String>,
}
#[derive(Debug, Clone)]
pub enum AgentImage {
Reference(String),
Full(DeployImage),
}
impl From<&str> for AgentImage {
fn from(reference: &str) -> Self {
AgentImage::Reference(reference.to_string())
}
}
impl From<String> for AgentImage {
fn from(reference: String) -> Self {
AgentImage::Reference(reference)
}
}
impl From<DeployImage> for AgentImage {
fn from(image: DeployImage) -> Self {
AgentImage::Full(image)
}
}
#[derive(Debug, Clone, Default)]
pub struct AgentScaling {
pub min: Option<u32>,
pub max: Option<u32>,
}
#[derive(Debug, Clone, Default)]
pub struct DeployAgentParams {
pub name: String,
pub image: Option<AgentImage>,
pub env: HashMap<String, String>,
pub scaling: Option<AgentScaling>,
pub profile: Option<AgentComputeProfile>,
pub region: Option<String>,
pub agent_id: Option<String>,
pub version_tag: Option<String>,
pub webhook: Option<RoomWebhook>,
pub env_secret_name: Option<String>,
}
#[derive(Debug, Clone)]
pub struct DeployedAgent {
pub deployment: AgentDeployment,
pub version: AgentDeploymentVersion,
pub secret_id: Option<String>,
}
fn normalize_deploy_image(image: Option<AgentImage>) -> DeployImage {
let mut image = match image {
None => DeployImage::default(),
Some(AgentImage::Reference(reference)) => DeployImage {
image_url: reference,
..Default::default()
},
Some(AgentImage::Full(image)) => image,
};
if image.image_cr.is_none() {
let first_segment = image.image_url.split('/').next().unwrap_or("");
image.image_cr = Some(
if !first_segment.is_empty() && first_segment.contains(['.', ':']) {
first_segment.to_string()
} else {
"docker.io".to_string()
},
);
}
if image.image_type.is_none() {
image.image_type = Some(
if image.image_pull_secret.is_some() {
"private"
} else {
"public"
}
.to_string(),
);
}
image
}
#[derive(Debug, Clone, Default)]
pub struct ListDeploymentsParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub agent_id: Option<String>,
pub name: Option<String>,
pub status: Option<String>,
pub deployment_id: Option<String>,
}
#[derive(Debug, Clone, Default)]
pub struct ListDeploymentLogsParams {
pub page: Option<u32>,
pub per_page: Option<u32>,
pub cursor: Option<String>,
pub agent_id: Option<String>,
pub version_id: Option<String>,
pub session_id: Option<String>,
pub user_id: Option<String>,
pub room_id: Option<String>,
pub start: Option<String>,
pub end: Option<String>,
pub sort: 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!(ListDeploymentsParams);
pagination!(ListDeploymentLogsParams);
#[derive(Debug, Clone, Copy)]
pub struct AgentDeploymentResource<'a> {
client: &'a Client,
}
impl<'a> AgentDeploymentResource<'a> {
pub async fn list(&self, params: ListDeploymentsParams) -> Result<Page<AgentDeployment>> {
paginate(
self.fetcher(¶ms),
¶ms.pagination(),
"deployments",
None,
)
.await
}
pub fn list_stream(
&self,
params: ListDeploymentsParams,
) -> impl Stream<Item = Result<AgentDeployment>> + Send {
auto_page(
self.fetcher(¶ms),
params.pagination(),
"deployments",
None,
)
}
pub async fn get(&self, deployment_id: &str) -> Result<AgentDeployment> {
let page = self
.list(ListDeploymentsParams {
deployment_id: Some(deployment_id.to_string()),
per_page: Some(1),
..Default::default()
})
.await?;
page.data
.iter()
.find(|deployment| deployment.id == deployment_id)
.or_else(|| page.data.first())
.cloned()
.ok_or_else(|| Error::not_found(format!("agent deployment {deployment_id} not found")))
}
pub async fn get_latest_version(&self, deployment_id: &str) -> Result<AgentDeploymentVersion> {
let path = format!("{DEPLOYMENTS_PATH}/{}/latest", escape(deployment_id));
self.client
.json(Method::GET, &path, CallOptions::new())
.await
}
pub async fn logs(
&self,
deployment_id: &str,
params: ListDeploymentLogsParams,
) -> Result<Page<AgentLogEntry>> {
let fetcher = self.logs_fetcher(deployment_id, ¶ms);
paginate(fetcher, ¶ms.pagination(), "logs", None).await
}
pub fn logs_stream(
&self,
deployment_id: &str,
params: ListDeploymentLogsParams,
) -> impl Stream<Item = Result<AgentLogEntry>> + Send {
let fetcher = self.logs_fetcher(deployment_id, ¶ms);
auto_page(fetcher, params.pagination(), "logs", None)
}
pub async fn delete(&self, deployment_id: &str, force: bool) -> Result<MessageResponse> {
let path = format!("{DEPLOYMENTS_PATH}/{}/delete", escape(deployment_id));
let body = json!({ "force": force });
self.client
.json(Method::POST, &path, CallOptions::json(&body)?)
.await
}
pub async fn create(&self, params: DeployAgentParams) -> Result<DeployedAgent> {
let image = normalize_deploy_image(params.image);
let mut secret_id = None;
if !params.env.is_empty() {
let name = params.env_secret_name.clone().unwrap_or_else(|| {
let suffix: String = random_uuid().chars().take(8).collect();
format!("{}-env-{suffix}", slugify_name(¶ms.name))
});
let body = json!({"name": name, "keys": params.env, "type": "NORMAL"});
#[derive(Deserialize)]
struct Secret {
id: String,
}
let secret: Secret = self
.client
.data(Method::POST, SECRETS_PATH, CallOptions::json(&body)?)
.await?;
secret_id = Some(secret.id);
}
let mut deploy_body = Map::new();
deploy_body.insert("name".into(), json!(params.name));
if let Some(agent_id) = params.agent_id.as_deref().filter(|id| !id.is_empty()) {
deploy_body.insert("agentId".into(), json!(agent_id));
}
if let Some(secret_id) = &secret_id {
deploy_body.insert("envSecret".into(), json!(secret_id));
}
let deployment: AgentDeployment = self
.client
.json(
Method::POST,
DEPLOYMENTS_PATH,
CallOptions::json(&deploy_body)?,
)
.await?;
let scaling = params.scaling.unwrap_or_default();
let version_body = VersionWire {
agent_id: deployment.agent_id.as_deref(),
deployment_id: Some(&deployment.id),
image,
profile: params.profile.unwrap_or(AgentComputeProfile::CPU_SMALL),
min_replica: scaling.min,
max_replica: scaling.max,
region: params.region.as_deref(),
version_tag: params.version_tag.as_deref(),
webhook: params.webhook.as_ref(),
};
let version: AgentDeploymentVersion = self
.client
.json(
Method::POST,
VERSIONS_PATH,
CallOptions::json(&version_body)?,
)
.await?;
Ok(DeployedAgent {
deployment,
version,
secret_id,
})
}
pub async fn set_version_active(
&self,
version_id: &str,
activate: bool,
force: bool,
) -> Result<MessageResponse> {
let path = format!("{VERSIONS_PATH}/{}/activate", escape(version_id));
let body = json!({"activate": activate, "force": force});
self.client
.json(Method::PUT, &path, CallOptions::json(&body)?)
.await
}
fn fetcher(&self, params: &ListDeploymentsParams) -> 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("agentId", params.agent_id.as_deref())
.opt_str("name", params.name.as_deref())
.opt_str("status", params.status.as_deref())
.opt_str("deploymentId", params.deployment_id.as_deref())
.into_pairs();
client
.json::<Value>(
Method::GET,
DEPLOYMENTS_PATH,
CallOptions::new().query(query),
)
.await
})
})
}
fn logs_fetcher(&self, deployment_id: &str, params: &ListDeploymentLogsParams) -> PageFetcher {
let client = self.client.clone();
let params = params.clone();
let path = format!("{DEPLOYMENTS_PATH}/{}/console-logs", escape(deployment_id));
Arc::new(move |page, per_page| {
let client = client.clone();
let params = params.clone();
let path = path.clone();
Box::pin(async move {
let query = QueryBuilder::new()
.opt("page", page)
.opt("perPage", per_page)
.opt_str("agentId", params.agent_id.as_deref())
.opt_str("versionId", params.version_id.as_deref())
.opt_str("sessionId", params.session_id.as_deref())
.opt_str("userId", params.user_id.as_deref())
.opt_str("roomId", params.room_id.as_deref())
.opt_str("start", params.start.as_deref())
.opt_str("end", params.end.as_deref())
.opt_str("sort", params.sort.as_deref())
.into_pairs();
client
.json::<Value>(Method::GET, &path, CallOptions::new().query(query))
.await
})
})
}
}
#[derive(Debug, Serialize)]
#[serde(rename_all = "camelCase")]
struct VersionWire<'a> {
#[serde(skip_serializing_if = "Option::is_none")]
agent_id: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
deployment_id: Option<&'a str>,
image: DeployImage,
profile: AgentComputeProfile,
#[serde(skip_serializing_if = "Option::is_none")]
min_replica: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
max_replica: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
region: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
version_tag: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
webhook: Option<&'a RoomWebhook>,
}
#[derive(Debug, Clone, Copy)]
pub struct AgentsResource<'a> {
client: &'a Client,
}
impl<'a> AgentsResource<'a> {
pub(crate) fn new(client: &'a Client) -> Self {
Self { client }
}
pub fn deployment(&self) -> AgentDeploymentResource<'a> {
AgentDeploymentResource {
client: self.client,
}
}
pub async fn dispatch(&self, params: DispatchAgentParams) -> Result<DispatchResult> {
self.send("/v2/agent/dispatch", params).await
}
pub async fn connect(&self, params: DispatchAgentParams) -> Result<DispatchResult> {
self.send("/v2/agent/connect", params).await
}
pub async fn general_dispatch(
&self,
params: GeneralDispatchAgentParams,
) -> Result<DispatchResult> {
self.send("/v2/agent/general/dispatch", params.into()).await
}
pub async fn init_config(&self) -> Result<AgentInitConfig> {
self.client
.data(Method::POST, "/v2/agent/init-config", CallOptions::new())
.await
}
async fn send(&self, path: &str, params: DispatchAgentParams) -> Result<DispatchResult> {
let meeting_id = params
.meeting_id
.as_deref()
.or(params.room_id.as_deref())
.filter(|id| !id.is_empty())
.ok_or_else(|| Error::validation("agents.dispatch() requires meeting_id (or room_id)"))?
.to_string();
let room_id = params
.room_id
.filter(|id| !id.is_empty())
.unwrap_or_else(|| meeting_id.clone());
let body = DispatchWire {
agent_id: params.agent_id,
room_options: params.room_options,
sip_options: params.sip_options,
metadata: params.metadata,
version_id: params.version_id,
version_tag: params.version_tag,
wait: params.wait,
meeting_id,
room_id: Some(room_id),
};
self.client
.data(Method::POST, path, CallOptions::json(&body)?)
.await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn a_registry_host_is_parsed_out_of_the_image_url() {
let image = normalize_deploy_image(Some("registry.videosdk.live/acme/agent:1.4".into()));
assert_eq!(image.image_cr.as_deref(), Some("registry.videosdk.live"));
assert_eq!(image.image_type.as_deref(), Some("public"));
let image = normalize_deploy_image(Some("localhost:5000/agent".into()));
assert_eq!(image.image_cr.as_deref(), Some("localhost:5000"));
}
#[test]
fn a_short_reference_falls_back_to_docker_hub() {
for reference in ["acme/agent:1.4", "agent", "library/agent"] {
let image = normalize_deploy_image(Some(reference.into()));
assert_eq!(image.image_cr.as_deref(), Some("docker.io"), "{reference}");
}
}
#[test]
fn a_pull_secret_makes_the_image_private() {
let image = normalize_deploy_image(Some(
DeployImage {
image_url: "registry.example.com/agent".into(),
image_pull_secret: Some("secret-1".into()),
..Default::default()
}
.into(),
));
assert_eq!(image.image_type.as_deref(), Some("private"));
assert_eq!(image.image_url, "registry.example.com/agent");
}
#[test]
fn an_explicit_image_object_is_never_overwritten() {
let image = normalize_deploy_image(Some(
DeployImage {
image_url: "registry.example.com/agent".into(),
image_cr: Some("custom.registry".into()),
image_type: Some("managed".into()),
image_pull_secret: None,
}
.into(),
));
assert_eq!(image.image_url, "registry.example.com/agent");
assert_eq!(image.image_cr.as_deref(), Some("custom.registry"));
assert_eq!(image.image_type.as_deref(), Some("managed"));
}
}