use futures_util::stream::{self, Stream, StreamExt, TryStreamExt};
use reqwest::Method;
use serde::{Deserialize, Serialize};
use crate::client::Ghl;
use crate::contacts::ListMeta;
use crate::error::Result;
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[allow(missing_docs)] pub struct Opportunity {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub location_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub pipeline_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub pipeline_stage_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub status: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub contact_id: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub monetary_value: Option<f64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub assigned_to: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub created_at: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub updated_at: Option<String>,
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
#[allow(missing_docs)] pub struct CreateOpportunity {
pub location_id: String,
pub pipeline_id: String,
pub name: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub pipeline_stage_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub contact_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub monetary_value: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub assigned_to: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize)]
#[serde(rename_all = "camelCase")]
#[allow(missing_docs)] pub struct UpdateOpportunity {
#[serde(skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub pipeline_stage_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub status: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub monetary_value: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
pub assigned_to: Option<String>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[allow(missing_docs)] pub struct Pipeline {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub stages: Vec<PipelineStage>,
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
#[allow(missing_docs)] pub struct PipelineStage {
pub id: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub name: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub position: Option<i64>,
#[serde(flatten)]
pub extra: serde_json::Map<String, serde_json::Value>,
}
#[derive(Deserialize)]
struct OpportunityEnvelope {
opportunity: Opportunity,
}
#[derive(Deserialize)]
struct PipelineList {
#[serde(default)]
pipelines: Vec<Pipeline>,
}
#[derive(Debug, Clone, Deserialize)]
pub struct OpportunityPage {
#[serde(default)]
pub opportunities: Vec<Opportunity>,
#[serde(default)]
pub meta: Option<ListMeta>,
}
pub struct OpportunitiesService {
client: Ghl,
}
impl OpportunitiesService {
pub(crate) fn new(client: Ghl) -> Self {
Self { client }
}
pub async fn pipelines(&self, location_id: &str) -> Result<Vec<Pipeline>> {
let list: PipelineList = self
.client
.send(
Method::GET,
"/opportunities/pipelines",
&[("locationId".into(), location_id.to_owned())],
None::<&()>,
)
.await?;
Ok(list.pipelines)
}
pub async fn create(&self, opportunity: CreateOpportunity) -> Result<Opportunity> {
let envelope: OpportunityEnvelope = self
.client
.send(Method::POST, "/opportunities/", &[], Some(&opportunity))
.await?;
Ok(envelope.opportunity)
}
pub async fn get(&self, opportunity_id: &str) -> Result<Opportunity> {
let envelope: OpportunityEnvelope = self
.client
.send(
Method::GET,
&format!("/opportunities/{opportunity_id}"),
&[],
None::<&()>,
)
.await?;
Ok(envelope.opportunity)
}
pub async fn update(
&self,
opportunity_id: &str,
update: UpdateOpportunity,
) -> Result<Opportunity> {
let envelope: OpportunityEnvelope = self
.client
.send(
Method::PUT,
&format!("/opportunities/{opportunity_id}"),
&[],
Some(&update),
)
.await?;
Ok(envelope.opportunity)
}
pub async fn update_status(&self, opportunity_id: &str, status: &str) -> Result<()> {
let _: serde_json::Value = self
.client
.send(
Method::PUT,
&format!("/opportunities/{opportunity_id}/status"),
&[],
Some(&serde_json::json!({ "status": status })),
)
.await?;
Ok(())
}
pub async fn delete(&self, opportunity_id: &str) -> Result<()> {
let _: serde_json::Value = self
.client
.send(
Method::DELETE,
&format!("/opportunities/{opportunity_id}"),
&[],
None::<&()>,
)
.await?;
Ok(())
}
pub fn search(&self, location_id: impl Into<String>) -> SearchOpportunities {
SearchOpportunities {
client: self.client.clone(),
location_id: location_id.into(),
limit: 20,
query: None,
pipeline_id: None,
status: None,
start_after_id: None,
start_after: None,
}
}
}
#[derive(Clone)]
pub struct SearchOpportunities {
client: Ghl,
location_id: String,
limit: u32,
query: Option<String>,
pipeline_id: Option<String>,
status: Option<String>,
start_after_id: Option<String>,
start_after: Option<i64>,
}
impl SearchOpportunities {
pub fn limit(mut self, limit: u32) -> Self {
self.limit = limit.clamp(1, 100);
self
}
pub fn query(mut self, query: impl Into<String>) -> Self {
self.query = Some(query.into());
self
}
pub fn pipeline_id(mut self, pipeline_id: impl Into<String>) -> Self {
self.pipeline_id = Some(pipeline_id.into());
self
}
pub fn status(mut self, status: impl Into<String>) -> Self {
self.status = Some(status.into());
self
}
pub fn start_after_id(mut self, cursor: impl Into<String>) -> Self {
self.start_after_id = Some(cursor.into());
self
}
pub async fn page(&self) -> Result<OpportunityPage> {
let mut query: Vec<(String, String)> = vec![
("location_id".into(), self.location_id.clone()),
("limit".into(), self.limit.to_string()),
];
if let Some(q) = &self.query {
query.push(("q".into(), q.clone()));
}
if let Some(p) = &self.pipeline_id {
query.push(("pipeline_id".into(), p.clone()));
}
if let Some(s) = &self.status {
query.push(("status".into(), s.clone()));
}
if let Some(id) = &self.start_after_id {
query.push(("startAfterId".into(), id.clone()));
}
if let Some(ts) = self.start_after {
query.push(("startAfter".into(), ts.to_string()));
}
self.client
.send(Method::GET, "/opportunities/search", &query, None::<&()>)
.await
}
pub fn stream(self) -> impl Stream<Item = Result<Opportunity>> {
stream::try_unfold(Some(self), |state| async move {
let Some(mut request) = state else {
return Ok::<_, crate::Error>(None);
};
let page = request.page().await?;
let full_page = page.opportunities.len() as u32 >= request.limit;
let cursor = page.meta.as_ref().and_then(|m| m.start_after_id.clone());
let start_after = page.meta.as_ref().and_then(|m| m.start_after);
let next = match (full_page, cursor) {
(true, Some(cursor)) => {
request.start_after_id = Some(cursor);
request.start_after = start_after;
Some(request)
}
_ => None,
};
Ok(Some((
stream::iter(page.opportunities.into_iter().map(Ok)),
next,
)))
})
.try_flatten()
.boxed()
}
}