use tracing::info;
use url::Url;
use crate::core::auth::{DeviceAuthChallenge, GoogleAuth, TokenStorage};
use crate::core::error::{GrrError, Result};
use crate::core::http::{HttpCore, QueryParams, TransportInfo, join_url, parse_url};
use crate::core::runtime::detect_runtime_features;
use crate::core::{Page, paginate};
use crate::people::models::*;
const DEFAULT_BASE_URL: &str = "https://people.googleapis.com/v1/";
const PROBE_PATH: &str = "people/me?personFields=names";
const CONNECTIONS_PATH: &str = "people/me/connections";
const CONTACT_GROUPS_PATH: &str = "contactGroups";
const OTHER_CONTACTS_PATH: &str = "otherContacts";
const SEARCH_CONTACTS_PATH: &str = "./people:searchContacts";
const CREATE_CONTACT_PATH: &str = "./people:createContact";
const BATCH_GET_PATH: &str = "./people:batchGet";
const READ_PERSON_FIELDS: &str = "names,emailAddresses,phoneNumbers,organizations,photos";
const SEARCH_READ_MASK: &str = "names,emailAddresses,phoneNumbers";
const GROUP_FIELDS: &str = "clientData,groupType,memberCount,metadata,name";
const CREATE_GROUP_READ_FIELDS: &str = "clientData,groupType,metadata,name";
const OTHER_CONTACTS_READ_MASK: &str = "names,emailAddresses";
const COPY_OTHER_CONTACT_MASK: &str = "names,emailAddresses,phoneNumbers";
const PHOTO_PERSON_FIELDS: &str = "photos";
const READ_SOURCE_TYPE_CONTACT: &str = "READ_SOURCE_TYPE_CONTACT";
const MAX_PAGE_SIZE: u32 = 1000;
const MAX_BATCH_GET_RESOURCE_NAMES: usize = 200;
const MAX_CONTACT_GROUP_MEMBERS_MODIFY: usize = 1000;
#[derive(Debug, Clone)]
pub struct ListConnectionsOptions {
pub page_size: u32,
pub query_name: Option<String>,
pub max_results: usize,
}
impl Default for ListConnectionsOptions {
fn default() -> Self {
Self {
page_size: 100,
query_name: None,
max_results: 0,
}
}
}
pub struct PeopleClientBuilder {
auth: Option<GoogleAuth>,
base_url: Option<Url>,
}
impl PeopleClientBuilder {
pub fn new() -> Self {
Self {
auth: None,
base_url: None,
}
}
pub fn auth(mut self, auth: GoogleAuth) -> Self {
self.auth = Some(auth);
self
}
pub fn base_url(mut self, url: Url) -> Self {
self.base_url = Some(url);
self
}
pub async fn build(self) -> Result<PeopleClient> {
let auth = self
.auth
.ok_or_else(|| GrrError::Config("Auth is required".into()))?;
let has_base_override = self.base_url.is_some();
let base_url = match self.base_url {
Some(url) => url,
None => parse_url(DEFAULT_BASE_URL, "base")?,
};
let core = if has_base_override {
HttpCore::unprobed(auth, crate::core::http::build_http_client()?)
} else {
HttpCore::connect(auth, &base_url, PROBE_PATH).await?
};
PeopleClient::new(core, base_url).await
}
}
impl Default for PeopleClientBuilder {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct PeopleClient {
core: HttpCore,
base_url: Url,
}
impl PeopleClient {
pub async fn new(core: HttpCore, base_url: Url) -> Result<Self> {
let features = detect_runtime_features().await;
info!(
"PeopleClient initialized: http3=always, io_uring={}",
features.io_uring
);
Ok(Self { core, base_url })
}
pub fn core(&self) -> &HttpCore {
&self.core
}
pub fn transport_info(&self) -> &TransportInfo {
self.core.transport_info()
}
pub fn token_backend(&self) -> &'static str {
self.core.auth().token_backend()
}
pub async fn login(&self) -> Result<TokenStorage> {
self.core.auth().login().await
}
pub async fn request_device_code(&self) -> Result<DeviceAuthChallenge> {
self.core.auth().request_device_code().await
}
pub async fn poll_device_code(
&self,
challenge: &mut DeviceAuthChallenge,
) -> Result<Option<TokenStorage>> {
self.core.auth().poll_device_code(challenge).await
}
pub async fn access_token(&self) -> Result<String> {
self.core.auth().get_access_token().await
}
fn api_url(&self, path: &str) -> Result<Url> {
join_url(&self.base_url, path, "API")
}
pub async fn list_connections(&self, opts: ListConnectionsOptions) -> Result<Vec<Person>> {
let batch_size = opts.page_size.clamp(1, MAX_PAGE_SIZE).to_string();
let query_name = opts.query_name.as_deref();
let max = if opts.max_results == 0 {
None
} else {
Some(opts.max_results)
};
paginate(max, move |page_token| {
let batch_size = batch_size.clone();
async move {
let params = QueryParams::new()
.add("personFields", READ_PERSON_FIELDS)
.add("pageSize", batch_size.as_str())
.add_optional("queryName", query_name)
.add_page_token(page_token.as_deref());
let page: ListConnectionsResponse = self
.core
.execute_json(params.apply(self.core.get(self.api_url(CONNECTIONS_PATH)?)))
.await?;
Ok(Page::new(
page.connections.unwrap_or_default(),
page.next_page_token,
))
}
})
.await
}
pub async fn search_contacts(&self, query: &str) -> Result<Vec<Person>> {
let request = self
.core
.get(self.api_url(SEARCH_CONTACTS_PATH)?)
.query(&[("query", query), ("readMask", SEARCH_READ_MASK)]);
let response: SearchContactsResponse = self.core.execute(request).await?.json().await?;
Ok(response
.results
.unwrap_or_default()
.into_iter()
.filter_map(|result| result.person)
.collect())
}
pub async fn get_person(&self, resource_name: &str) -> Result<Person> {
let request = self
.core
.get(self.api_url(resource_name)?)
.query(&[("personFields", READ_PERSON_FIELDS)]);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn create_contact(&self, person: Person) -> Result<Person> {
let request = self
.core
.post(self.api_url(CREATE_CONTACT_PATH)?)
.json(&person);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn update_contact(
&self,
resource_name: &str,
person: Person,
update_fields: &[&str],
) -> Result<Person> {
if update_fields.is_empty() {
return Err(GrrError::InvalidArgument(
"update_contact: at least one update field is required".into(),
));
}
let mask = update_fields.join(",");
let path = format!("{resource_name}:updateContact");
let request = self
.core
.patch(self.api_url(&path)?)
.query(&[("updatePersonFields", mask.as_str())])
.json(&person);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn delete_contact(&self, resource_name: &str) -> Result<()> {
let path = format!("{resource_name}:deleteContact");
self.core
.execute(self.core.delete(self.api_url(&path)?))
.await?;
Ok(())
}
pub async fn list_contact_groups(&self) -> Result<Vec<ContactGroup>> {
let page_size = MAX_PAGE_SIZE.to_string();
paginate(None, move |page_token| {
let page_size = page_size.clone();
async move {
let params = QueryParams::new()
.add("groupFields", GROUP_FIELDS)
.add("pageSize", page_size.as_str())
.add_page_token(page_token.as_deref());
let page: ListContactGroupsResponse = self
.core
.execute_json(params.apply(self.core.get(self.api_url(CONTACT_GROUPS_PATH)?)))
.await?;
Ok(Page::new(
page.contact_groups.unwrap_or_default(),
page.next_page_token,
))
}
})
.await
}
pub async fn get_contact_group(&self, resource_name: &str) -> Result<ContactGroup> {
validate_resource_name(resource_name, "contactGroups")?;
let request = self
.core
.get(self.api_url(resource_name)?)
.query(&[("groupFields", GROUP_FIELDS)]);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn create_contact_group(&self, name: &str) -> Result<ContactGroup> {
validate_nonempty(name, "contact group name")?;
let body = CreateContactGroupRequest {
contact_group: ContactGroup {
name: Some(name.to_owned()),
..ContactGroup::default()
},
read_group_fields: Some(CREATE_GROUP_READ_FIELDS.to_owned()),
};
let request = self
.core
.post(self.api_url(CONTACT_GROUPS_PATH)?)
.json(&body);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn update_contact_group(
&self,
resource_name: &str,
name: &str,
) -> Result<ContactGroup> {
validate_resource_name(resource_name, "contactGroups")?;
validate_nonempty(name, "contact group name")?;
let body = UpdateContactGroupRequest {
contact_group: ContactGroup {
resource_name: Some(resource_name.to_owned()),
name: Some(name.to_owned()),
..ContactGroup::default()
},
update_group_fields: Some("name".to_owned()),
read_group_fields: Some(GROUP_FIELDS.to_owned()),
};
let request = self.core.put(self.api_url(resource_name)?).json(&body);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn delete_contact_group(&self, resource_name: &str) -> Result<()> {
validate_resource_name(resource_name, "contactGroups")?;
self.core
.execute(self.core.delete(self.api_url(resource_name)?))
.await?;
Ok(())
}
pub async fn modify_contact_group_members(
&self,
resource_name: &str,
resource_names_to_add: &[String],
resource_names_to_remove: &[String],
) -> Result<ModifyContactGroupMembersResponse> {
validate_resource_name(resource_name, "contactGroups")?;
let total = resource_names_to_add
.len()
.checked_add(resource_names_to_remove.len())
.ok_or_else(|| GrrError::InvalidArgument("too many contact group members".into()))?;
if total == 0 {
return Err(GrrError::InvalidArgument(
"at least one contact group member is required".into(),
));
}
if total > MAX_CONTACT_GROUP_MEMBERS_MODIFY {
return Err(GrrError::InvalidArgument(format!(
"contact group member limit is {MAX_CONTACT_GROUP_MEMBERS_MODIFY}"
)));
}
for resource_name in resource_names_to_add.iter().chain(resource_names_to_remove) {
validate_resource_name(resource_name, "people")?;
}
let body = ModifyContactGroupMembersRequest {
resource_names_to_add: resource_names_to_add.to_vec(),
resource_names_to_remove: resource_names_to_remove.to_vec(),
};
let path = format!("{resource_name}/members:modify");
let request = self.core.post(self.api_url(&path)?).json(&body);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn get_people(
&self,
resource_names: &[String],
person_fields: &str,
) -> Result<GetPeopleResponse> {
validate_nonempty(person_fields, "personFields")?;
if resource_names.is_empty() {
return Err(GrrError::InvalidArgument(
"at least one resource name is required".into(),
));
}
if resource_names.len() > MAX_BATCH_GET_RESOURCE_NAMES {
return Err(GrrError::InvalidArgument(format!(
"batch get supports at most {MAX_BATCH_GET_RESOURCE_NAMES} resource names"
)));
}
for resource_name in resource_names {
validate_resource_name(resource_name, "people")?;
}
let mut params = QueryParams::new();
for resource_name in resource_names {
params = params.add("resourceNames", resource_name);
}
let request = params
.add("personFields", person_fields)
.apply(self.core.get(self.api_url(BATCH_GET_PATH)?));
self.core.execute_json(request).await
}
pub async fn list_other_contacts(&self, max_results: usize) -> Result<Vec<Person>> {
let max = if max_results == 0 {
None
} else {
Some(max_results)
};
let page_size = if max_results == 0 {
MAX_PAGE_SIZE
} else {
u32::try_from(max_results.min(MAX_PAGE_SIZE as usize)).unwrap_or(MAX_PAGE_SIZE)
}
.to_string();
paginate(max, move |page_token| {
let page_size = page_size.clone();
async move {
let params = QueryParams::new()
.add("readMask", OTHER_CONTACTS_READ_MASK)
.add("pageSize", page_size.as_str())
.add_page_token(page_token.as_deref());
let page: ListOtherContactsResponse = self
.core
.execute_json(params.apply(self.core.get(self.api_url(OTHER_CONTACTS_PATH)?)))
.await?;
Ok(Page::new(
page.other_contacts.unwrap_or_default(),
page.next_page_token,
))
}
})
.await
}
pub async fn copy_other_contact_to_my_contacts_group(
&self,
resource_name: &str,
) -> Result<Person> {
validate_resource_name(resource_name, "otherContacts")?;
let body = CopyOtherContactToMyContactsGroupRequest {
copy_mask: COPY_OTHER_CONTACT_MASK.to_owned(),
read_mask: None,
sources: Some(vec![READ_SOURCE_TYPE_CONTACT.to_owned()]),
};
let path = format!("{resource_name}:copyOtherContactToMyContactsGroup");
let request = self.core.post(self.api_url(&path)?).json(&body);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn list_photos(&self, resource_name: &str) -> Result<Vec<Photo>> {
validate_resource_name(resource_name, "people")?;
let request = self
.core
.get(self.api_url(resource_name)?)
.query(&[("personFields", PHOTO_PERSON_FIELDS)]);
let person: Person = self.core.execute(request).await?.json().await?;
Ok(person.photos.unwrap_or_default())
}
pub async fn update_contact_photo(
&self,
resource_name: &str,
photo_bytes: &str,
) -> Result<UpdateContactPhotoResponse> {
validate_resource_name(resource_name, "people")?;
validate_nonempty(photo_bytes, "base64 photo bytes")?;
let body = UpdateContactPhotoRequest {
photo_bytes: photo_bytes.to_owned(),
person_fields: Some(PHOTO_PERSON_FIELDS.to_owned()),
sources: Some(vec![READ_SOURCE_TYPE_CONTACT.to_owned()]),
};
let path = format!("{resource_name}:updateContactPhoto");
let request = self.core.patch(self.api_url(&path)?).json(&body);
Ok(self.core.execute(request).await?.json().await?)
}
pub async fn delete_contact_photo(
&self,
resource_name: &str,
) -> Result<DeleteContactPhotoResponse> {
validate_resource_name(resource_name, "people")?;
let path = format!("{resource_name}:deleteContactPhoto");
let request = self
.core
.delete(self.api_url(&path)?)
.query(&[("personFields", PHOTO_PERSON_FIELDS)]);
Ok(self.core.execute(request).await?.json().await?)
}
}
fn validate_resource_name(resource_name: &str, collection: &str) -> Result<()> {
let prefix = format!("{collection}/");
let valid = resource_name
.strip_prefix(&prefix)
.and_then(|id| (!id.is_empty() && !id.contains('/')).then_some(id))
.is_some();
if !valid {
return Err(GrrError::InvalidArgument(format!(
"resource name must have the form {collection}/<id>"
)));
}
Ok(())
}
fn validate_nonempty(value: &str, field: &str) -> Result<()> {
if value.trim().is_empty() {
return Err(GrrError::InvalidArgument(format!("{field} is required")));
}
Ok(())
}