1use std::sync::atomic::{AtomicI64, Ordering};
4use std::sync::Arc;
5use std::time::Duration;
6
7use reqwest::header::{HeaderMap, HeaderValue, AUTHORIZATION};
8use reqwest::{Method, StatusCode};
9use serde::de::DeserializeOwned;
10use serde::Serialize;
11
12use crate::auth::Auth;
13use crate::calendars::CalendarsService;
14use crate::contacts::ContactsService;
15use crate::conversations::ConversationsService;
16use crate::error::{Error, Result};
17use crate::locations::LocationsService;
18use crate::opportunities::OpportunitiesService;
19
20pub const DEFAULT_BASE_URL: &str = "https://services.leadconnectorhq.com";
22
23pub const API_VERSION: &str = "2021-07-28";
25
26const DEFAULT_TIMEOUT: Duration = Duration::from_secs(30);
27const DEFAULT_MAX_RETRIES: u32 = 3;
28const BACKOFF_BASE: Duration = Duration::from_millis(500);
29const BACKOFF_CAP: Duration = Duration::from_secs(8);
30
31#[derive(Clone)]
46pub struct Ghl {
47 inner: Arc<Inner>,
48}
49
50struct Inner {
51 http: reqwest::Client,
52 base_url: String,
53 auth: Auth,
54 max_retries: u32,
55 rate_remaining: AtomicI64,
57 rate_daily_remaining: AtomicI64,
59}
60
61impl std::fmt::Debug for Ghl {
62 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
63 f.debug_struct("Ghl")
64 .field("base_url", &self.inner.base_url)
65 .field("auth", &self.inner.auth)
66 .finish_non_exhaustive()
67 }
68}
69
70#[derive(Default)]
72pub struct GhlBuilder {
73 base_url: Option<String>,
74 auth: Option<Auth>,
75 timeout: Option<Duration>,
76 max_retries: Option<u32>,
77}
78
79impl GhlBuilder {
80 pub fn base_url(mut self, url: impl Into<String>) -> Self {
82 self.base_url = Some(url.into());
83 self
84 }
85
86 pub fn private_integration_token(mut self, token: impl Into<String>) -> Self {
88 self.auth = Some(Auth::private_integration(token));
89 self
90 }
91
92 pub fn access_token(mut self, token: impl Into<String>) -> Self {
94 self.auth = Some(Auth::access_token(token));
95 self
96 }
97
98 pub fn auth(mut self, auth: Auth) -> Self {
100 self.auth = Some(auth);
101 self
102 }
103
104 pub fn timeout(mut self, timeout: Duration) -> Self {
106 self.timeout = Some(timeout);
107 self
108 }
109
110 pub fn max_retries(mut self, retries: u32) -> Self {
112 self.max_retries = Some(retries);
113 self
114 }
115
116 pub fn build(self) -> Result<Ghl> {
118 let auth = self.auth.ok_or_else(|| {
119 Error::Config(
120 "no credentials configured — call `.private_integration_token(…)`, \
121 `.access_token(…)`, or `.auth(…)`"
122 .into(),
123 )
124 })?;
125 let base_url = self
126 .base_url
127 .unwrap_or_else(|| DEFAULT_BASE_URL.to_owned())
128 .trim_end_matches('/')
129 .to_owned();
130 let http = reqwest::Client::builder()
131 .timeout(self.timeout.unwrap_or(DEFAULT_TIMEOUT))
132 .user_agent(concat!("ghl-sdk/", env!("CARGO_PKG_VERSION")))
133 .build()?;
134 Ok(Ghl {
135 inner: Arc::new(Inner {
136 http,
137 base_url,
138 auth,
139 max_retries: self.max_retries.unwrap_or(DEFAULT_MAX_RETRIES),
140 rate_remaining: AtomicI64::new(-1),
141 rate_daily_remaining: AtomicI64::new(-1),
142 }),
143 })
144 }
145}
146
147#[derive(Debug, Clone, Copy, PartialEq, Eq)]
149pub struct RateStatus {
150 pub burst_remaining: Option<i64>,
152 pub daily_remaining: Option<i64>,
154}
155
156impl Ghl {
157 pub fn builder() -> GhlBuilder {
159 GhlBuilder::default()
160 }
161
162 pub fn from_env() -> Result<Self> {
168 let mut builder = Ghl::builder();
169 if let Ok(url) = std::env::var("GHL_BASE_URL") {
170 builder = builder.base_url(url);
171 }
172 if let Ok(token) = std::env::var("GHL_PIT_TOKEN") {
173 builder = builder.private_integration_token(token);
174 } else if let Ok(token) = std::env::var("GHL_ACCESS_TOKEN") {
175 builder = builder.access_token(token);
176 } else {
177 return Err(Error::Config(
178 "set GHL_PIT_TOKEN (or GHL_ACCESS_TOKEN) in the environment, \
179 or use `Ghl::builder()` to pass credentials as parameters"
180 .into(),
181 ));
182 }
183 builder.build()
184 }
185
186 pub fn contacts(&self) -> ContactsService {
188 ContactsService::new(self.clone())
189 }
190
191 pub fn locations(&self) -> LocationsService {
193 LocationsService::new(self.clone())
194 }
195
196 pub fn opportunities(&self) -> OpportunitiesService {
198 OpportunitiesService::new(self.clone())
199 }
200
201 pub fn conversations(&self) -> ConversationsService {
203 ConversationsService::new(self.clone())
204 }
205
206 pub fn calendars(&self) -> CalendarsService {
208 CalendarsService::new(self.clone())
209 }
210
211 #[cfg(feature = "ad-manager")]
215 #[cfg_attr(docsrs, doc(cfg(feature = "ad-manager")))]
216 pub fn ad_manager(&self) -> crate::services::ad_manager::AdManagerService {
217 crate::services::ad_manager::AdManagerService::new(self.clone())
218 }
219
220 #[cfg(feature = "affiliate-manager")]
224 #[cfg_attr(docsrs, doc(cfg(feature = "affiliate-manager")))]
225 pub fn affiliate_manager(&self) -> crate::services::affiliate_manager::AffiliateManagerService {
226 crate::services::affiliate_manager::AffiliateManagerService::new(self.clone())
227 }
228
229 #[cfg(feature = "agent-studio")]
233 #[cfg_attr(docsrs, doc(cfg(feature = "agent-studio")))]
234 pub fn agent_studio(&self) -> crate::services::agent_studio::AgentStudioService {
235 crate::services::agent_studio::AgentStudioService::new(self.clone())
236 }
237
238 #[cfg(feature = "associations")]
242 #[cfg_attr(docsrs, doc(cfg(feature = "associations")))]
243 pub fn associations(&self) -> crate::services::associations::AssociationsService {
244 crate::services::associations::AssociationsService::new(self.clone())
245 }
246
247 #[cfg(feature = "blogs")]
251 #[cfg_attr(docsrs, doc(cfg(feature = "blogs")))]
252 pub fn blogs(&self) -> crate::services::blogs::BlogsService {
253 crate::services::blogs::BlogsService::new(self.clone())
254 }
255
256 #[cfg(feature = "brand-boards")]
260 #[cfg_attr(docsrs, doc(cfg(feature = "brand-boards")))]
261 pub fn brand_boards(&self) -> crate::services::brand_boards::BrandBoardsService {
262 crate::services::brand_boards::BrandBoardsService::new(self.clone())
263 }
264
265 #[cfg(feature = "businesses")]
269 #[cfg_attr(docsrs, doc(cfg(feature = "businesses")))]
270 pub fn businesses(&self) -> crate::services::businesses::BusinessesService {
271 crate::services::businesses::BusinessesService::new(self.clone())
272 }
273
274 #[cfg(feature = "campaigns")]
278 #[cfg_attr(docsrs, doc(cfg(feature = "campaigns")))]
279 pub fn campaigns(&self) -> crate::services::campaigns::CampaignsService {
280 crate::services::campaigns::CampaignsService::new(self.clone())
281 }
282
283 #[cfg(feature = "companies")]
287 #[cfg_attr(docsrs, doc(cfg(feature = "companies")))]
288 pub fn companies(&self) -> crate::services::companies::CompaniesService {
289 crate::services::companies::CompaniesService::new(self.clone())
290 }
291
292 #[cfg(feature = "conversation-ai")]
296 #[cfg_attr(docsrs, doc(cfg(feature = "conversation-ai")))]
297 pub fn conversation_ai(&self) -> crate::services::conversation_ai::ConversationAiService {
298 crate::services::conversation_ai::ConversationAiService::new(self.clone())
299 }
300
301 #[cfg(feature = "courses")]
305 #[cfg_attr(docsrs, doc(cfg(feature = "courses")))]
306 pub fn courses(&self) -> crate::services::courses::CoursesService {
307 crate::services::courses::CoursesService::new(self.clone())
308 }
309
310 #[cfg(feature = "custom-fields")]
314 #[cfg_attr(docsrs, doc(cfg(feature = "custom-fields")))]
315 pub fn custom_fields(&self) -> crate::services::custom_fields::CustomFieldsService {
316 crate::services::custom_fields::CustomFieldsService::new(self.clone())
317 }
318
319 #[cfg(feature = "custom-menus")]
323 #[cfg_attr(docsrs, doc(cfg(feature = "custom-menus")))]
324 pub fn custom_menus(&self) -> crate::services::custom_menus::CustomMenusService {
325 crate::services::custom_menus::CustomMenusService::new(self.clone())
326 }
327
328 #[cfg(feature = "email-isv")]
332 #[cfg_attr(docsrs, doc(cfg(feature = "email-isv")))]
333 pub fn email_isv(&self) -> crate::services::email_isv::EmailIsvService {
334 crate::services::email_isv::EmailIsvService::new(self.clone())
335 }
336
337 #[cfg(feature = "emails")]
341 #[cfg_attr(docsrs, doc(cfg(feature = "emails")))]
342 pub fn emails(&self) -> crate::services::emails::EmailsService {
343 crate::services::emails::EmailsService::new(self.clone())
344 }
345
346 #[cfg(feature = "forms")]
350 #[cfg_attr(docsrs, doc(cfg(feature = "forms")))]
351 pub fn forms(&self) -> crate::services::forms::FormsService {
352 crate::services::forms::FormsService::new(self.clone())
353 }
354
355 #[cfg(feature = "funnels")]
359 #[cfg_attr(docsrs, doc(cfg(feature = "funnels")))]
360 pub fn funnels(&self) -> crate::services::funnels::FunnelsService {
361 crate::services::funnels::FunnelsService::new(self.clone())
362 }
363
364 #[cfg(feature = "invoices")]
368 #[cfg_attr(docsrs, doc(cfg(feature = "invoices")))]
369 pub fn invoices(&self) -> crate::services::invoices::InvoicesService {
370 crate::services::invoices::InvoicesService::new(self.clone())
371 }
372
373 #[cfg(feature = "knowledge-base")]
377 #[cfg_attr(docsrs, doc(cfg(feature = "knowledge-base")))]
378 pub fn knowledge_base(&self) -> crate::services::knowledge_base::KnowledgeBaseService {
379 crate::services::knowledge_base::KnowledgeBaseService::new(self.clone())
380 }
381
382 #[cfg(feature = "links")]
386 #[cfg_attr(docsrs, doc(cfg(feature = "links")))]
387 pub fn links(&self) -> crate::services::links::LinksService {
388 crate::services::links::LinksService::new(self.clone())
389 }
390
391 #[cfg(feature = "marketplace")]
395 #[cfg_attr(docsrs, doc(cfg(feature = "marketplace")))]
396 pub fn marketplace(&self) -> crate::services::marketplace::MarketplaceService {
397 crate::services::marketplace::MarketplaceService::new(self.clone())
398 }
399
400 #[cfg(feature = "medias")]
404 #[cfg_attr(docsrs, doc(cfg(feature = "medias")))]
405 pub fn medias(&self) -> crate::services::medias::MediasService {
406 crate::services::medias::MediasService::new(self.clone())
407 }
408
409 #[cfg(feature = "oauth")]
413 #[cfg_attr(docsrs, doc(cfg(feature = "oauth")))]
414 pub fn oauth(&self) -> crate::services::oauth::OauthService {
415 crate::services::oauth::OauthService::new(self.clone())
416 }
417
418 #[cfg(feature = "objects")]
422 #[cfg_attr(docsrs, doc(cfg(feature = "objects")))]
423 pub fn objects(&self) -> crate::services::objects::ObjectsService {
424 crate::services::objects::ObjectsService::new(self.clone())
425 }
426
427 #[cfg(feature = "payments")]
431 #[cfg_attr(docsrs, doc(cfg(feature = "payments")))]
432 pub fn payments(&self) -> crate::services::payments::PaymentsService {
433 crate::services::payments::PaymentsService::new(self.clone())
434 }
435
436 #[cfg(feature = "phone-system")]
440 #[cfg_attr(docsrs, doc(cfg(feature = "phone-system")))]
441 pub fn phone_system(&self) -> crate::services::phone_system::PhoneSystemService {
442 crate::services::phone_system::PhoneSystemService::new(self.clone())
443 }
444
445 #[cfg(feature = "products")]
449 #[cfg_attr(docsrs, doc(cfg(feature = "products")))]
450 pub fn products(&self) -> crate::services::products::ProductsService {
451 crate::services::products::ProductsService::new(self.clone())
452 }
453
454 #[cfg(feature = "proposals")]
458 #[cfg_attr(docsrs, doc(cfg(feature = "proposals")))]
459 pub fn proposals(&self) -> crate::services::proposals::ProposalsService {
460 crate::services::proposals::ProposalsService::new(self.clone())
461 }
462
463 #[cfg(feature = "saas-api")]
467 #[cfg_attr(docsrs, doc(cfg(feature = "saas-api")))]
468 pub fn saas_api(&self) -> crate::services::saas_api::SaasApiService {
469 crate::services::saas_api::SaasApiService::new(self.clone())
470 }
471
472 #[cfg(feature = "snapshots")]
476 #[cfg_attr(docsrs, doc(cfg(feature = "snapshots")))]
477 pub fn snapshots(&self) -> crate::services::snapshots::SnapshotsService {
478 crate::services::snapshots::SnapshotsService::new(self.clone())
479 }
480
481 #[cfg(feature = "social-media-posting")]
485 #[cfg_attr(docsrs, doc(cfg(feature = "social-media-posting")))]
486 pub fn social_media_posting(
487 &self,
488 ) -> crate::services::social_media_posting::SocialMediaPostingService {
489 crate::services::social_media_posting::SocialMediaPostingService::new(self.clone())
490 }
491
492 #[cfg(feature = "store")]
496 #[cfg_attr(docsrs, doc(cfg(feature = "store")))]
497 pub fn store(&self) -> crate::services::store::StoreService {
498 crate::services::store::StoreService::new(self.clone())
499 }
500
501 #[cfg(feature = "surveys")]
505 #[cfg_attr(docsrs, doc(cfg(feature = "surveys")))]
506 pub fn surveys(&self) -> crate::services::surveys::SurveysService {
507 crate::services::surveys::SurveysService::new(self.clone())
508 }
509
510 #[cfg(feature = "users")]
514 #[cfg_attr(docsrs, doc(cfg(feature = "users")))]
515 pub fn users(&self) -> crate::services::users::UsersService {
516 crate::services::users::UsersService::new(self.clone())
517 }
518
519 #[cfg(feature = "voice-ai")]
523 #[cfg_attr(docsrs, doc(cfg(feature = "voice-ai")))]
524 pub fn voice_ai(&self) -> crate::services::voice_ai::VoiceAiService {
525 crate::services::voice_ai::VoiceAiService::new(self.clone())
526 }
527
528 #[cfg(feature = "workflows")]
532 #[cfg_attr(docsrs, doc(cfg(feature = "workflows")))]
533 pub fn workflows(&self) -> crate::services::workflows::WorkflowsService {
534 crate::services::workflows::WorkflowsService::new(self.clone())
535 }
536
537 pub fn v3(&self) -> crate::services::v3::V3 {
547 crate::services::v3::V3 {
548 client: self.clone(),
549 }
550 }
551 pub fn rate_status(&self) -> RateStatus {
553 let read = |a: &AtomicI64| {
554 let v = a.load(Ordering::Relaxed);
555 (v >= 0).then_some(v)
556 };
557 RateStatus {
558 burst_remaining: read(&self.inner.rate_remaining),
559 daily_remaining: read(&self.inner.rate_daily_remaining),
560 }
561 }
562
563 pub async fn as_location(&self, company_id: &str, location_id: &str) -> Result<Ghl> {
568 #[derive(serde::Deserialize)]
569 struct LocationTokenResponse {
570 access_token: String,
571 }
572
573 let bearer = self
574 .inner
575 .auth
576 .bearer(&self.inner.http, &self.inner.base_url)
577 .await?;
578 let response = self
579 .inner
580 .http
581 .post(format!("{}/oauth/locationToken", self.inner.base_url))
582 .header(AUTHORIZATION, format!("Bearer {bearer}"))
583 .header("Version", API_VERSION)
584 .form(&[("companyId", company_id), ("locationId", location_id)])
585 .send()
586 .await?;
587
588 let status = response.status();
589 if !status.is_success() {
590 let body = response.text().await.unwrap_or_default();
591 return Err(Error::Auth(format!(
592 "location token exchange failed ({status}): {body}"
593 )));
594 }
595 let parsed: LocationTokenResponse = response
596 .json()
597 .await
598 .map_err(|e| Error::Auth(format!("unexpected locationToken response: {e}")))?;
599
600 Ghl::builder()
601 .base_url(&self.inner.base_url)
602 .access_token(parsed.access_token)
603 .max_retries(self.inner.max_retries)
604 .build()
605 }
606
607 pub async fn get_raw(&self, path: &str, query: &[(&str, &str)]) -> Result<serde_json::Value> {
611 let query: Vec<(String, String)> = query
612 .iter()
613 .map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
614 .collect();
615 self.send(Method::GET, path, &query, None::<&()>).await
616 }
617
618 pub async fn post_raw(&self, path: &str, body: &impl Serialize) -> Result<serde_json::Value> {
620 self.send(Method::POST, path, &[], Some(body)).await
621 }
622
623 pub async fn request_raw(
630 &self,
631 method: &str,
632 path: &str,
633 query: &[(String, String)],
634 body: Option<&serde_json::Value>,
635 version: Option<&str>,
636 ) -> Result<serde_json::Value> {
637 let method = Method::from_bytes(method.to_uppercase().as_bytes())
638 .map_err(|_| Error::Config(format!("invalid HTTP method `{method}`")))?;
639 self.send_versioned(method, path, query, body, version)
640 .await
641 }
642
643 pub(crate) async fn send<T: DeserializeOwned>(
646 &self,
647 method: Method,
648 path: &str,
649 query: &[(String, String)],
650 body: Option<&impl Serialize>,
651 ) -> Result<T> {
652 self.send_versioned(method, path, query, body, None).await
653 }
654
655 pub(crate) async fn send_versioned<T: DeserializeOwned>(
656 &self,
657 method: Method,
658 path: &str,
659 query: &[(String, String)],
660 body: Option<&impl Serialize>,
661 version: Option<&str>,
662 ) -> Result<T> {
663 let body = match body {
665 Some(b) => Some(serde_json::to_value(b).map_err(|source| Error::Decode {
666 endpoint: path.to_owned(),
667 source,
668 })?),
669 None => None,
670 };
671
672 let url = format!("{}{}", self.inner.base_url, path);
673 let idempotent = matches!(
674 method,
675 Method::GET | Method::PUT | Method::DELETE | Method::HEAD
676 );
677 let mut attempt: u32 = 0;
678
679 loop {
680 let bearer = self
681 .inner
682 .auth
683 .bearer(&self.inner.http, &self.inner.base_url)
684 .await?;
685
686 let mut request = self
687 .inner
688 .http
689 .request(method.clone(), &url)
690 .header(AUTHORIZATION, format!("Bearer {bearer}"))
691 .header("Version", version.unwrap_or(API_VERSION))
692 .header(reqwest::header::ACCEPT, "application/json");
693 if !query.is_empty() {
694 request = request.query(query);
695 }
696 if let Some(ref b) = body {
697 request = request.json(b);
698 }
699
700 let outcome = request.send().await;
701
702 match outcome {
703 Ok(response) => {
704 self.record_rate_headers(response.headers());
705 let status = response.status();
706
707 if status.is_success() {
708 let bytes = response.bytes().await?;
709 return serde_json::from_slice(&bytes).map_err(|source| Error::Decode {
710 endpoint: path.to_owned(),
711 source,
712 });
713 }
714
715 let retry_after = parse_retry_after(response.headers());
716 let request_id = response
717 .headers()
718 .get("x-request-id")
719 .and_then(|v| v.to_str().ok())
720 .map(str::to_owned);
721 let message = read_api_message(response).await;
722
723 let retryable = status == StatusCode::TOO_MANY_REQUESTS
724 || (idempotent && status.is_server_error());
725 if retryable && attempt < self.inner.max_retries {
726 let delay = retry_after.unwrap_or_else(|| backoff_delay(attempt));
727 tracing::warn!(
728 %status, attempt, delay_ms = delay.as_millis() as u64, path,
729 "GoHighLevel request failed; retrying"
730 );
731 tokio::time::sleep(delay).await;
732 attempt += 1;
733 continue;
734 }
735
736 return Err(if status == StatusCode::TOO_MANY_REQUESTS {
737 Error::RateLimited { retry_after }
738 } else {
739 Error::Api {
740 status,
741 message,
742 request_id,
743 }
744 });
745 }
746 Err(err) => {
747 if idempotent && attempt < self.inner.max_retries && err.status().is_none() {
749 let delay = backoff_delay(attempt);
750 tracing::warn!(
751 error = %err, attempt, delay_ms = delay.as_millis() as u64, path,
752 "transport error; retrying"
753 );
754 tokio::time::sleep(delay).await;
755 attempt += 1;
756 continue;
757 }
758 return Err(err.into());
759 }
760 }
761 }
762 }
763
764 fn record_rate_headers(&self, headers: &HeaderMap) {
765 let parse =
766 |name: &str| -> Option<i64> { headers.get(name)?.to_str().ok()?.trim().parse().ok() };
767 if let Some(v) = parse("x-ratelimit-remaining") {
768 self.inner.rate_remaining.store(v, Ordering::Relaxed);
769 }
770 if let Some(v) = parse("x-ratelimit-daily-remaining") {
771 self.inner.rate_daily_remaining.store(v, Ordering::Relaxed);
772 }
773 }
774}
775
776fn backoff_delay(attempt: u32) -> Duration {
778 let exp = BACKOFF_BASE.saturating_mul(2u32.saturating_pow(attempt));
779 let cap = exp.min(BACKOFF_CAP);
780 cap.mul_f64(0.5 + fastrand::f64() * 0.5)
781}
782
783fn parse_retry_after(headers: &HeaderMap) -> Option<Duration> {
784 let value: &HeaderValue = headers.get(reqwest::header::RETRY_AFTER)?;
785 let seconds: u64 = value.to_str().ok()?.trim().parse().ok()?;
786 Some(Duration::from_secs(seconds))
787}
788
789async fn read_api_message(response: reqwest::Response) -> String {
792 let text = response.text().await.unwrap_or_default();
793 match serde_json::from_str::<serde_json::Value>(&text) {
794 Ok(v) => match v.get("message") {
795 Some(serde_json::Value::String(s)) => s.clone(),
796 Some(serde_json::Value::Array(parts)) => parts
797 .iter()
798 .filter_map(|p| p.as_str())
799 .collect::<Vec<_>>()
800 .join("; "),
801 _ => text,
802 },
803 Err(_) => text,
804 }
805}