Skip to main content

powerplatform_dataverse_client/dataverse/
serviceclient.rs

1use std::collections::HashMap;
2use std::path::PathBuf;
3
4use chrono::{DateTime, Utc};
5use log::debug;
6use reqwest::Client;
7use reqwest::header::CONTENT_TYPE;
8use serde::de::DeserializeOwned;
9use serde_json::Map;
10use serde_json::Value;
11use tokio::sync::Mutex;
12use uuid::Uuid;
13
14use crate::LogLevel;
15use crate::auth::config::AuthConfig;
16use crate::auth::connectionstring::{
17    parse_connection_string_auth_config, parse_connection_string_url,
18};
19use crate::auth::credentials::{TokenExchange, refresh_device_code_token};
20use crate::auth::token::{
21    CachedToken, fetch_token_for_config, is_expiring_soon, load_cached_token,
22    resolve_token_cache_file_path, save_cached_token,
23};
24use crate::dataverse::batch::{
25    ExecuteMultipleRequest, ExecuteMultipleResponse, ExecuteMultipleResponseItem,
26    OrganizationRequest, ParsedBatchPart, PreparedBatchItem, PreparedBatchRequest,
27    entity_to_write_body, parse_batch_response_parts, parse_fault,
28};
29use crate::dataverse::entity::Entity;
30use crate::dataverse::entity::Value::Int;
31use crate::dataverse::entityattribute::EntityAttribute;
32use crate::dataverse::entitydefinition::EntityDefinition;
33use crate::dataverse::entityrelationship::EntityRelationship;
34use crate::dataverse::fetchxml::{apply_paging, ensure_aggregate_page_size, fetch_tag_has_attr};
35use crate::dataverse::parse::{
36    extract_paging_cookie, parse_entities_from_response, parse_more_records,
37    parse_record_count_from_response,
38};
39use crate::dataverse::requestparameters::RequestParameters;
40
41const ROW_NUMBER_ATTRIBUTE: &str = "__rownum";
42const AGGREGATE_PAGE_SIZE: i32 = 5000;
43const DEFAULT_FETCHXML_PAGE_SIZE: i32 = 5000;
44
45/// OData list wrapper returned by Dataverse metadata endpoints.
46#[derive(Debug, serde::Deserialize)]
47struct ODataList<T> {
48    value: Vec<T>,
49}
50
51#[derive(Debug, serde::Deserialize)]
52struct EntityRelationshipDirectional {
53    #[serde(rename = "SchemaName")]
54    schema_name: String,
55    #[serde(rename = "ReferencedEntity")]
56    referenced_entity: Option<String>,
57    #[serde(rename = "ReferencedAttribute")]
58    referenced_attribute: Option<String>,
59    #[serde(rename = "ReferencingEntity")]
60    referencing_entity: Option<String>,
61    #[serde(rename = "ReferencingAttribute")]
62    referencing_attribute: Option<String>,
63    #[serde(rename = "IsCustomRelationship")]
64    is_custom_relationship: Option<bool>,
65    #[serde(flatten)]
66    extra: Map<String, Value>,
67}
68
69#[derive(Debug, serde::Deserialize)]
70struct EntityRelationshipManyToMany {
71    #[serde(rename = "SchemaName")]
72    schema_name: String,
73    #[serde(rename = "Entity1LogicalName")]
74    entity1_logical_name: Option<String>,
75    #[serde(rename = "Entity2LogicalName")]
76    entity2_logical_name: Option<String>,
77    #[serde(rename = "IntersectEntityName")]
78    intersect_entity_name: Option<String>,
79    #[serde(rename = "IsCustomRelationship")]
80    is_custom_relationship: Option<bool>,
81    #[serde(flatten)]
82    extra: Map<String, Value>,
83}
84
85/// HTTP client for Dataverse Web API operations.
86pub struct ServiceClient {
87    client: Client,
88    auth: AuthConfig,
89    base_url: std::string::String,
90    token_cache_path: PathBuf,
91    token: Mutex<CachedToken>,
92    // Entity definitions are cached as a single blob because most metadata-driven features need
93    // the full list, and Dataverse returns them efficiently in one request.
94    entity_definitions_cache: Mutex<Option<Vec<EntityDefinition>>>,
95    // Attribute metadata is cached per logical entity name because callers usually fan out to only
96    // a small number of entities during a session.
97    entity_attributes_cache: Mutex<HashMap<String, Vec<EntityAttribute>>>,
98    log_level: LogLevel,
99}
100
101impl ServiceClient {
102    /// Create a new client from a Dataverse connection string.
103    pub async fn new(connection_string: &str, log_level: LogLevel) -> Result<Self, String> {
104        let base_url = parse_connection_string_url(connection_string)?;
105        let auth = parse_connection_string_auth_config(connection_string)?;
106        Self::new_internal(auth, base_url, log_level).await
107    }
108
109    /// Create a new client from explicit authentication configuration.
110    pub async fn new_with_auth(auth: AuthConfig, log_level: LogLevel) -> Result<Self, String> {
111        let base_url = auth.dataverse_url().to_string();
112        Self::new_internal(auth, base_url, log_level).await
113    }
114
115    async fn new_internal(
116        auth: AuthConfig,
117        base_url: String,
118        log_level: LogLevel,
119    ) -> Result<Self, String> {
120        let token_cache_path = resolve_token_cache_file_path(&auth)?;
121
122        // Initialization eagerly ensures a usable token so later requests can fail on Dataverse
123        // semantics instead of first-request authentication setup.
124        let token = if let Some(cached) = load_cached_token(&token_cache_path)? {
125            if !cached.access_token.trim().is_empty() && !is_expiring_soon(cached.expires_at) {
126                cached
127            } else {
128                let refreshed = fetch_token_for_config(&auth).await?;
129                save_cached_token(&token_cache_path, &refreshed)?;
130                refreshed
131            }
132        } else {
133            let fetched = fetch_token_for_config(&auth).await?;
134            save_cached_token(&token_cache_path, &fetched)?;
135            fetched
136        };
137
138        Ok(Self {
139            client: Client::new(),
140            auth,
141            base_url,
142            token_cache_path,
143            token: Mutex::new(token),
144            entity_definitions_cache: Mutex::new(None),
145            entity_attributes_cache: Mutex::new(HashMap::new()),
146            log_level,
147        })
148    }
149
150    /// Return the current token expiry as a UTC datetime.
151    pub async fn token_expires_at(&self) -> Option<DateTime<Utc>> {
152        let expires_at = self.token.lock().await.expires_at?;
153        DateTime::<Utc>::from_timestamp(expires_at as i64, 0)
154    }
155
156    /// Retrieve a single FetchXML response page without automatic paging.
157    pub async fn retrieve_multiple_fetchxml(
158        &self,
159        entity: &str,
160        fetchxml: &str,
161    ) -> Result<Vec<Entity>, std::string::String> {
162        let primary_id_attribute = self.resolve_primary_id_attribute(entity).await?;
163        let attribute_map = self.entity_attribute_map(entity).await?;
164        self.retrieve_multiple_fetchxml_single(
165            entity,
166            fetchxml,
167            primary_id_attribute.as_deref(),
168            Some(&attribute_map),
169        )
170        .await
171    }
172
173    /// Retrieve multiple records by FetchXML, automatically paging until all results are returned.
174    pub async fn retrieve_multiple_fetchxml_paging(
175        &self,
176        entity: &str,
177        fetchxml: &str,
178    ) -> Result<Vec<Entity>, std::string::String> {
179        self.retrieve_multiple_fetchxml_paging_with_progress(entity, fetchxml, |_, _| {}, None)
180            .await
181    }
182
183    /// Retrieve multiple records by FetchXML, automatically paging until all results are returned.
184    /// Uses the provided page size when specified, otherwise defaults to 5000 records per page.
185    /// Reports page-level progress as `(page_number, total_records_retrieved_so_far)`.
186    pub async fn retrieve_multiple_fetchxml_paging_with_progress<F>(
187        &self,
188        entity: &str,
189        fetchxml: &str,
190        mut on_progress: F,
191        page_size: Option<i32>,
192    ) -> Result<Vec<Entity>, std::string::String>
193    where
194        F: FnMut(usize, usize),
195    {
196        let page_size = page_size.unwrap_or(DEFAULT_FETCHXML_PAGE_SIZE);
197        let primary_id_attribute = self.resolve_primary_id_attribute(entity).await?;
198        let attribute_map = self.entity_attribute_map(entity).await?;
199        if fetch_tag_has_attr(fetchxml, "top")? {
200            let entities = self
201                .retrieve_multiple_fetchxml_single(
202                    entity,
203                    fetchxml,
204                    primary_id_attribute.as_deref(),
205                    Some(&attribute_map),
206                )
207                .await?;
208            on_progress(1, entities.len());
209            return Ok(entities);
210        }
211
212        let mut page = 1;
213        let mut paging_cookie: Option<std::string::String> = None;
214        let mut entities: Vec<Entity> = vec![];
215
216        loop {
217            let fetchxml = ensure_fetch_page_size(fetchxml, page_size)?;
218            let fetch_with_paging = apply_paging(
219                &ensure_aggregate_page_size(&fetchxml, AGGREGATE_PAGE_SIZE)?,
220                page,
221                paging_cookie.as_deref(),
222            )?;
223
224            if self.log_level.includes_debug() {
225                debug!("Fetch page: {}", page);
226                debug!("FetchXML: {}", fetch_with_paging);
227            }
228
229            let mut url = format!("{}/api/data/v9.2/{}", self.base_url, entity);
230            url.push_str("?fetchXml=");
231            url.push_str(&urlencoding::encode(&fetch_with_paging));
232
233            if self.log_level.includes_debug() {
234                debug!("Url: {:?}", url);
235            }
236
237            let access_token = self.get_access_token().await?;
238            let resp = self
239                .client
240                .get(&url)
241                .bearer_auth(&access_token)
242                .header("Accept", "application/json")
243                .header(
244                    "Prefer",
245                    "odata.include-annotations=\"Microsoft.Dynamics.CRM.fetchxmlpagingcookie,Microsoft.Dynamics.CRM.morerecords,Microsoft.Dynamics.CRM.lookuplogicalname,OData.Community.Display.V1.FormattedValue\"",
246                )
247                .send()
248                .await
249                .map_err(|e| format!("Request failed: {e}"))?;
250
251            let status = resp.status();
252
253            if !status.is_success() {
254                let body = resp.text().await.unwrap_or_default();
255                return Err(format!("Dataverse API error ({}): {}", status, body));
256            }
257
258            let json: Value = resp
259                .json()
260                .await
261                .map_err(|e| format!("Failed to parse JSON: {e}"))?;
262
263            let mut page_entities = parse_entities_from_response(
264                &json,
265                entity,
266                primary_id_attribute.as_deref(),
267                Some(&attribute_map),
268            )?;
269            let start_index = entities.len();
270            for (offset, entity) in page_entities.iter_mut().enumerate() {
271                let row_number = (start_index + offset + 1) as i64;
272                entity
273                    .attributes
274                    .insert(ROW_NUMBER_ATTRIBUTE.to_string(), Int(row_number));
275            }
276            entities.extend(page_entities);
277            on_progress(page as usize, entities.len());
278
279            let more_records = parse_more_records(&json);
280            if !more_records {
281                break;
282            }
283
284            paging_cookie = extract_paging_cookie(&json);
285            page += 1;
286        }
287
288        Ok(entities)
289    }
290
291    /// Count records for a FetchXML query without retrieving all data.
292    pub async fn retrieve_multiple_fetchxml_count(
293        &self,
294        entity: &str,
295        fetchxml: &str,
296    ) -> Result<usize, std::string::String> {
297        if fetch_tag_has_attr(fetchxml, "top")? {
298            let resp = self
299                .retrieve_multiple_fetchxml_single(entity, fetchxml, None, None)
300                .await?;
301            return Ok(resp.len());
302        }
303
304        let mut page = 1;
305        let mut paging_cookie: Option<std::string::String> = None;
306        let mut total = 0usize;
307
308        loop {
309            let fetch_with_paging = apply_paging(
310                &ensure_aggregate_page_size(fetchxml, AGGREGATE_PAGE_SIZE)?,
311                page,
312                paging_cookie.as_deref(),
313            )?;
314
315            if self.log_level.includes_debug() {
316                debug!("Fetch page: {}", page);
317                debug!("FetchXML: {}", fetch_with_paging);
318            }
319
320            let mut url = format!("{}/api/data/v9.2/{}", self.base_url, entity);
321            url.push_str("?fetchXml=");
322            url.push_str(&urlencoding::encode(&fetch_with_paging));
323
324            if self.log_level.includes_debug() {
325                debug!("Url: {:?}", url);
326            }
327
328            let access_token = self.get_access_token().await?;
329            let resp = self
330                .client
331                .get(&url)
332                .bearer_auth(&access_token)
333                .header("Accept", "application/json")
334                .header(
335                    "Prefer",
336                    "odata.include-annotations=\"Microsoft.Dynamics.CRM.fetchxmlpagingcookie,Microsoft.Dynamics.CRM.morerecords,Microsoft.Dynamics.CRM.lookuplogicalname,OData.Community.Display.V1.FormattedValue\"",
337                )
338                .send()
339                .await
340                .map_err(|e| format!("Request failed: {e}"))?;
341
342            let status = resp.status();
343
344            if !status.is_success() {
345                let body = resp.text().await.unwrap_or_default();
346                return Err(format!("Dataverse API error ({}): {}", status, body));
347            }
348
349            let json: Value = resp
350                .json()
351                .await
352                .map_err(|e| format!("Failed to parse JSON: {e}"))?;
353
354            total += parse_record_count_from_response(&json)?;
355
356            let more_records = parse_more_records(&json);
357            if !more_records {
358                break;
359            }
360
361            paging_cookie = extract_paging_cookie(&json);
362            page += 1;
363        }
364
365        Ok(total)
366    }
367
368    /// Retrieve a single page of FetchXML results.
369    async fn retrieve_multiple_fetchxml_single(
370        &self,
371        entity: &str,
372        fetchxml: &str,
373        primary_id_attribute: Option<&str>,
374        entity_attributes: Option<&HashMap<String, EntityAttribute>>,
375    ) -> Result<Vec<Entity>, std::string::String> {
376        if self.log_level.includes_debug() {
377            debug!("FetchXML: {}", fetchxml);
378        }
379
380        let mut url = format!("{}/api/data/v9.2/{}", self.base_url, entity);
381        url.push_str("?fetchXml=");
382        url.push_str(&urlencoding::encode(fetchxml));
383
384        if self.log_level.includes_debug() {
385            debug!("Url: {:?}", url);
386        }
387
388        let access_token = self.get_access_token().await?;
389        let resp = self
390            .client
391            .get(&url)
392            .bearer_auth(&access_token)
393            .header("Accept", "application/json")
394            .header(
395                "Prefer",
396                "odata.include-annotations=\"Microsoft.Dynamics.CRM.fetchxmlpagingcookie,Microsoft.Dynamics.CRM.morerecords,Microsoft.Dynamics.CRM.lookuplogicalname,OData.Community.Display.V1.FormattedValue\"",
397            )
398            .send()
399            .await
400            .map_err(|e| format!("Request failed: {e}"))?;
401
402        let status = resp.status();
403
404        if !status.is_success() {
405            let body = resp.text().await.unwrap_or_default();
406            return Err(format!("Dataverse API error ({}): {}", status, body));
407        }
408
409        let json: Value = resp
410            .json()
411            .await
412            .map_err(|e| format!("Failed to parse JSON: {e}"))?;
413
414        parse_entities_from_response(&json, entity, primary_id_attribute, entity_attributes)
415    }
416
417    /// List all entity definitions.
418    pub async fn list_entity_definitions(
419        &self,
420    ) -> Result<Vec<EntityDefinition>, std::string::String> {
421        {
422            let cache = self.entity_definitions_cache.lock().await;
423            if let Some(value) = &*cache {
424                return Ok(value.clone());
425            }
426        }
427
428        let url = format!(
429            "{}/api/data/v9.2/EntityDefinitions?$select=LogicalName,SchemaName,DisplayName,EntitySetName,IsCustomEntity,IsActivity,PrimaryIdAttribute",
430            self.base_url
431        );
432
433        let access_token = self.get_access_token().await?;
434        let resp = self
435            .client
436            .get(&url)
437            .bearer_auth(&access_token)
438            .header("Accept", "application/json")
439            .send()
440            .await
441            .map_err(|e| format!("Request failed: {e}"))?;
442
443        let status = resp.status();
444
445        if !status.is_success() {
446            let body = resp.text().await.unwrap_or_default();
447            return Err(format!("Dataverse API error ({}): {}", status, body));
448        }
449
450        let parsed: ODataList<EntityDefinition> = resp
451            .json()
452            .await
453            .map_err(|e| format!("Failed to parse JSON: {e}"))?;
454
455        let value = parsed.value;
456        let mut cache = self.entity_definitions_cache.lock().await;
457        *cache = Some(value.clone());
458
459        Ok(value)
460    }
461
462    /// List entity attributes for a given logical name.
463    pub async fn list_entity_attributes(
464        &self,
465        logical_name: &str,
466    ) -> Result<Vec<EntityAttribute>, std::string::String> {
467        {
468            let cache = self.entity_attributes_cache.lock().await;
469            if let Some(value) = cache.get(&normalize_entity_name(logical_name)) {
470                return Ok(value.clone());
471            }
472        }
473
474        let logical = logical_name.replace('\'', "''");
475        let url = format!(
476            "{}/api/data/v9.2/EntityDefinitions(LogicalName='{}')/Attributes?$select=LogicalName,SchemaName,AttributeType,AttributeTypeName,IsCustomAttribute,IsValidODataAttribute,IsValidForRead,IsValidForUpdate&$filter=IsValidODataAttribute eq true and IsValidForRead eq true",
477            self.base_url, logical
478        );
479
480        let access_token = self.get_access_token().await?;
481        let resp = self
482            .client
483            .get(&url)
484            .bearer_auth(&access_token)
485            .header("Accept", "application/json")
486            .send()
487            .await
488            .map_err(|e| format!("Request failed: {e}"))?;
489
490        let status = resp.status();
491
492        if !status.is_success() {
493            let body = resp.text().await.unwrap_or_default();
494            return Err(format!("Dataverse API error ({}): {}", status, body));
495        }
496
497        let parsed: ODataList<EntityAttribute> = resp
498            .json()
499            .await
500            .map_err(|e| format!("Failed to parse JSON: {e}"))?;
501
502        let value = parsed.value;
503        let mut cache = self.entity_attributes_cache.lock().await;
504        cache.insert(normalize_entity_name(logical_name), value.clone());
505
506        Ok(value)
507    }
508
509    /// List entity relationships for a given logical name.
510    pub async fn list_entity_relationships(
511        &self,
512        logical_name: &str,
513    ) -> Result<Vec<EntityRelationship>, std::string::String> {
514        let logical = logical_name.replace('\'', "''");
515        let many_to_one = self
516            .list_metadata_collection::<EntityRelationshipDirectional>(&format!(
517                "EntityDefinitions(LogicalName='{}')/ManyToOneRelationships?$select=SchemaName,ReferencedEntity,ReferencedAttribute,ReferencingEntity,ReferencingAttribute,IsCustomRelationship",
518                logical
519            ))
520            .await?
521            .into_iter()
522            .map(|relationship| EntityRelationship {
523                schema_name: relationship.schema_name,
524                relationship_type: "ManyToOne".to_string(),
525                referenced_entity: relationship.referenced_entity,
526                referenced_attribute: relationship.referenced_attribute,
527                referencing_entity: relationship.referencing_entity,
528                referencing_attribute: relationship.referencing_attribute,
529                intersect_entity_name: None,
530                is_custom_relationship: relationship.is_custom_relationship,
531                extra: relationship.extra.into_iter().collect(),
532            });
533
534        let one_to_many = self
535            .list_metadata_collection::<EntityRelationshipDirectional>(&format!(
536                "EntityDefinitions(LogicalName='{}')/OneToManyRelationships?$select=SchemaName,ReferencedEntity,ReferencedAttribute,ReferencingEntity,ReferencingAttribute,IsCustomRelationship",
537                logical
538            ))
539            .await?
540            .into_iter()
541            .map(|relationship| EntityRelationship {
542                schema_name: relationship.schema_name,
543                relationship_type: "OneToMany".to_string(),
544                referenced_entity: relationship.referenced_entity,
545                referenced_attribute: relationship.referenced_attribute,
546                referencing_entity: relationship.referencing_entity,
547                referencing_attribute: relationship.referencing_attribute,
548                intersect_entity_name: None,
549                is_custom_relationship: relationship.is_custom_relationship,
550                extra: relationship.extra.into_iter().collect(),
551            });
552
553        let many_to_many = self
554            .list_metadata_collection::<EntityRelationshipManyToMany>(&format!(
555                "EntityDefinitions(LogicalName='{}')/ManyToManyRelationships?$select=SchemaName,Entity1LogicalName,Entity2LogicalName,IntersectEntityName,IsCustomRelationship",
556                logical
557            ))
558            .await?
559            .into_iter()
560            .map(|relationship| EntityRelationship {
561                schema_name: relationship.schema_name,
562                relationship_type: "ManyToMany".to_string(),
563                referenced_entity: relationship.entity1_logical_name,
564                referenced_attribute: None,
565                referencing_entity: relationship.entity2_logical_name,
566                referencing_attribute: None,
567                intersect_entity_name: relationship.intersect_entity_name,
568                is_custom_relationship: relationship.is_custom_relationship,
569                extra: relationship.extra.into_iter().collect(),
570            });
571
572        Ok(many_to_one.chain(one_to_many).chain(many_to_many).collect())
573    }
574
575    /// Update a single entity record by ID.
576    pub async fn update_entity(
577        &self,
578        entity_set: &str,
579        id: &str,
580        attributes: &HashMap<std::string::String, Value>,
581    ) -> Result<(), std::string::String> {
582        self.update_entity_with_options(entity_set, id, attributes, &RequestParameters::default())
583            .await
584    }
585
586    /// Create a single entity record and return its ID when available.
587    pub async fn create_entity(
588        &self,
589        entity_set: &str,
590        attributes: &HashMap<std::string::String, Value>,
591    ) -> Result<Option<Uuid>, std::string::String> {
592        self.create_entity_with_options(entity_set, attributes, &RequestParameters::default())
593            .await
594    }
595
596    /// Create a single entity record with Dataverse request parameters.
597    pub async fn create_entity_with_options(
598        &self,
599        entity_set: &str,
600        attributes: &HashMap<std::string::String, Value>,
601        options: &RequestParameters,
602    ) -> Result<Option<Uuid>, std::string::String> {
603        let url = format!("{}/api/data/v9.2/{}", self.base_url, entity_set);
604
605        let access_token = self.get_access_token().await?;
606        let request = self
607            .client
608            .post(&url)
609            .bearer_auth(&access_token)
610            .header("Accept", "application/json")
611            .header("Content-Type", "application/json")
612            .json(attributes);
613
614        let resp = options
615            .apply(request)
616            .send()
617            .await
618            .map_err(|e| format!("Request failed: {e}"))?;
619
620        let status = resp.status();
621        if !status.is_success() {
622            let body = resp.text().await.unwrap_or_default();
623            return Err(format!("Dataverse API error ({}): {}", status, body));
624        }
625
626        Ok(resp
627            .headers()
628            .get("OData-EntityId")
629            .or_else(|| resp.headers().get("Location"))
630            .and_then(|value| value.to_str().ok())
631            .and_then(parse_uuid_from_uri))
632    }
633
634    /// Update a single entity record by ID with Dataverse request parameters.
635    pub async fn update_entity_with_options(
636        &self,
637        entity_set: &str,
638        id: &str,
639        attributes: &HashMap<std::string::String, Value>,
640        options: &RequestParameters,
641    ) -> Result<(), std::string::String> {
642        let trimmed = id.trim_matches(|ch| ch == '{' || ch == '}');
643        let url = format!(
644            "{}/api/data/v9.2/{}({})",
645            self.base_url, entity_set, trimmed
646        );
647
648        let access_token = self.get_access_token().await?;
649        let request = self
650            .client
651            .patch(&url)
652            .bearer_auth(&access_token)
653            .header("Accept", "application/json")
654            .header("Content-Type", "application/json")
655            .json(&attributes);
656
657        let resp = options
658            .apply(request)
659            .send()
660            .await
661            .map_err(|e| format!("Request failed: {e}"))?;
662
663        let status = resp.status();
664        if !status.is_success() {
665            let body = resp.text().await.unwrap_or_default();
666            return Err(format!("Dataverse API error ({}): {}", status, body));
667        }
668
669        Ok(())
670    }
671
672    /// Delete a single entity record by ID.
673    pub async fn delete_entity(
674        &self,
675        entity_set: &str,
676        id: &str,
677    ) -> Result<(), std::string::String> {
678        self.delete_entity_with_options(entity_set, id, &RequestParameters::default())
679            .await
680    }
681
682    /// Delete a single entity record by ID with Dataverse request parameters.
683    pub async fn delete_entity_with_options(
684        &self,
685        entity_set: &str,
686        id: &str,
687        options: &RequestParameters,
688    ) -> Result<(), std::string::String> {
689        let trimmed = id.trim_matches(|ch| ch == '{' || ch == '}');
690        let url = format!(
691            "{}/api/data/v9.2/{}({})",
692            self.base_url, entity_set, trimmed
693        );
694
695        let access_token = self.get_access_token().await?;
696        let request = self
697            .client
698            .delete(&url)
699            .bearer_auth(&access_token)
700            .header("Accept", "application/json");
701
702        let resp = options
703            .apply(request)
704            .send()
705            .await
706            .map_err(|e| format!("Request failed: {e}"))?;
707
708        let status = resp.status();
709        if !status.is_success() {
710            let body = resp.text().await.unwrap_or_default();
711            return Err(format!("Dataverse API error ({}): {}", status, body));
712        }
713
714        Ok(())
715    }
716
717    /// Execute multiple create, update, and delete requests using a single Dataverse batch call.
718    pub async fn execute_multiple(
719        &self,
720        request: &ExecuteMultipleRequest,
721    ) -> Result<ExecuteMultipleResponse, String> {
722        if request.requests.is_empty() {
723            return Ok(ExecuteMultipleResponse::default());
724        }
725
726        if request.requests.len() > 1000 {
727            return Err(format!(
728                "ExecuteMultipleRequest contains {} requests, exceeding the Dataverse batch limit of 1000",
729                request.requests.len()
730            ));
731        }
732
733        let entity_set_name_by_logical_name = self.entity_set_name_map().await?;
734        let prepared_requests = self
735            .prepare_batch_requests(&request.requests, &entity_set_name_by_logical_name)?;
736        let boundary = format!("batch_{}", Uuid::new_v4().as_hyphenated());
737        let body = self.build_batch_body(&boundary, &prepared_requests);
738        let url = format!("{}/api/data/v9.2/$batch", self.base_url);
739        let access_token = self.get_access_token().await?;
740
741        let mut http_request = self
742            .client
743            .post(&url)
744            .bearer_auth(&access_token)
745            .header("OData-MaxVersion", "4.0")
746            .header("OData-Version", "4.0")
747            .header("If-None-Match", "null")
748            .header("Accept", "application/json")
749            .header("Content-Type", format!("multipart/mixed; boundary={boundary}"))
750            .body(body);
751
752        if request.settings.continue_on_error {
753            http_request = http_request.header("Prefer", "odata.continue-on-error");
754        }
755
756        let resp = http_request
757            .send()
758            .await
759            .map_err(|e| format!("Request failed: {e}"))?;
760
761        let status = resp.status();
762        let content_type = resp
763            .headers()
764            .get(CONTENT_TYPE)
765            .and_then(|value| value.to_str().ok())
766            .map(|value| value.to_string());
767        let response_text = resp
768            .text()
769            .await
770            .map_err(|e| format!("Failed to read batch response: {e}"))?;
771
772        if !status.is_success() && !content_type.as_deref().unwrap_or_default().starts_with("multipart/mixed") {
773            return Err(format!(
774                "Dataverse API error ({}): {}",
775                status,
776                response_text
777            ));
778        }
779
780        let parts = parse_batch_response_parts(content_type.as_deref(), &response_text)?;
781        self.map_batch_response(request, parts)
782    }
783
784    async fn get_access_token(&self) -> Result<String, String> {
785        let mut token = self.token.lock().await;
786        if !token.access_token.trim().is_empty() && !is_expiring_soon(token.expires_at) {
787            return Ok(token.access_token.clone());
788        }
789
790        // Refreshing while the mutex is held keeps parallel callers from racing into multiple token
791        // refreshes and then stomping each other's cache file updates.
792        let refreshed = match &self.auth {
793            AuthConfig::ClientCredentials { .. } => {
794                println!("Refreshing access token before request using client credentials...");
795                fetch_token_for_config(&self.auth).await?
796            }
797            AuthConfig::DeviceCode {
798                client_id,
799                dataverse_url,
800                tenant_id,
801                ..
802            } => {
803                println!("Refreshing access token before request using device code...");
804                let refresh_token = token.refresh_token.clone().ok_or(
805                    "Device code token cannot refresh without a refresh token".to_string(),
806                )?;
807                let scope = format!(
808                    "{}/user_impersonation offline_access openid profile",
809                    dataverse_url
810                );
811                let token: TokenExchange =
812                    refresh_device_code_token(client_id, tenant_id, &scope, &refresh_token).await?;
813
814                CachedToken {
815                    access_token: token.access_token,
816                    refresh_token: Some(token.refresh_token),
817                    expires_at: Some(token.expires_at),
818                }
819            }
820        };
821
822        save_cached_token(&self.token_cache_path, &refreshed)?;
823        let access_token = refreshed.access_token.clone();
824        *token = refreshed;
825        Ok(access_token)
826    }
827
828    async fn list_metadata_collection<T>(&self, path: &str) -> Result<Vec<T>, String>
829    where
830        T: DeserializeOwned,
831    {
832        let url = format!("{}/api/data/v9.2/{}", self.base_url, path);
833
834        let access_token = self.get_access_token().await?;
835        let resp = self
836            .client
837            .get(&url)
838            .bearer_auth(&access_token)
839            .header("Accept", "application/json")
840            .send()
841            .await
842            .map_err(|e| format!("Request failed: {e}"))?;
843
844        let status = resp.status();
845
846        if !status.is_success() {
847            let body = resp.text().await.unwrap_or_default();
848            return Err(format!("Dataverse API error ({}): {}", status, body));
849        }
850
851        let parsed: ODataList<T> = resp
852            .json()
853            .await
854            .map_err(|e| format!("Failed to parse JSON: {e}"))?;
855
856        Ok(parsed.value)
857    }
858
859    async fn resolve_primary_id_attribute(
860        &self,
861        entity_set: &str,
862    ) -> Result<Option<String>, String> {
863        let definitions = self.list_entity_definitions().await?;
864        let target = normalize_entity_name(entity_set);
865
866        Ok(definitions
867            .into_iter()
868            .find(|definition| {
869                normalize_entity_name(&definition.entity_set_name) == target
870                    || normalize_entity_name(&definition.logical_name) == target
871                    || normalize_entity_name(&definition.schema_name) == target
872            })
873            .and_then(|definition| definition.primary_id_attribute))
874    }
875
876    async fn resolve_entity_logical_name(&self, entity_name: &str) -> Result<String, String> {
877        let definitions = self.list_entity_definitions().await?;
878        let target = normalize_entity_name(entity_name);
879
880        definitions
881            .into_iter()
882            .find(|definition| {
883                normalize_entity_name(&definition.entity_set_name) == target
884                    || normalize_entity_name(&definition.logical_name) == target
885                    || normalize_entity_name(&definition.schema_name) == target
886            })
887            .map(|definition| definition.logical_name)
888            .ok_or_else(|| format!("Entity metadata not found for '{}'", entity_name))
889    }
890
891    async fn entity_attribute_map(
892        &self,
893        entity_name: &str,
894    ) -> Result<HashMap<String, EntityAttribute>, String> {
895        let logical_name = self.resolve_entity_logical_name(entity_name).await?;
896        let attributes = self.list_entity_attributes(&logical_name).await?;
897        let mut map = HashMap::new();
898        for attribute in attributes {
899            map.insert(attribute.logical_name.to_ascii_lowercase(), attribute.clone());
900            map.entry(attribute.schema_name.to_ascii_lowercase())
901                .or_insert(attribute);
902        }
903        Ok(map)
904    }
905
906    async fn entity_set_name_map(&self) -> Result<HashMap<String, String>, String> {
907        let definitions = self.list_entity_definitions().await?;
908        Ok(definitions
909            .into_iter()
910            .map(|definition| {
911                (
912                    definition.logical_name.to_ascii_lowercase(),
913                    definition.entity_set_name,
914                )
915            })
916            .collect())
917    }
918
919    fn prepare_batch_requests(
920        &self,
921        requests: &[OrganizationRequest],
922        entity_set_name_by_logical_name: &HashMap<String, String>,
923    ) -> Result<Vec<PreparedBatchItem>, String> {
924        requests
925            .iter()
926            .enumerate()
927            .map(|(request_index, request)| {
928                self.prepare_batch_request(
929                    request_index,
930                    request.clone(),
931                    entity_set_name_by_logical_name,
932                )
933            })
934            .collect()
935    }
936
937    fn prepare_batch_request(
938        &self,
939        _request_index: usize,
940        request: OrganizationRequest,
941        entity_set_name_by_logical_name: &HashMap<String, String>,
942    ) -> Result<PreparedBatchItem, String> {
943        let prepared = match &request {
944            OrganizationRequest::Create(request) => {
945                let entity_set_name = entity_set_name_by_logical_name
946                    .get(&request.target.logical_name.to_ascii_lowercase())
947                    .ok_or_else(|| {
948                        format!(
949                            "Entity set metadata not found for '{}'",
950                            request.target.logical_name
951                        )
952                    })?;
953
954                PreparedBatchRequest {
955                    method: "POST",
956                    path: format!("/api/data/v9.2/{entity_set_name}"),
957                    body: Some(entity_to_write_body(
958                        &request.target,
959                        entity_set_name_by_logical_name,
960                    )?),
961                    parameters: request.parameters.clone(),
962                }
963            }
964            OrganizationRequest::Update(request) => {
965                if request.target.id.is_nil() {
966                    return Err("UpdateRequest target must include a non-empty entity ID".to_string());
967                }
968
969                let entity_set_name = entity_set_name_by_logical_name
970                    .get(&request.target.logical_name.to_ascii_lowercase())
971                    .ok_or_else(|| {
972                        format!(
973                            "Entity set metadata not found for '{}'",
974                            request.target.logical_name
975                        )
976                    })?;
977
978                PreparedBatchRequest {
979                    method: "PATCH",
980                    path: format!(
981                        "/api/data/v9.2/{}({})",
982                        entity_set_name,
983                        request.target.id.as_hyphenated()
984                    ),
985                    body: Some(entity_to_write_body(
986                        &request.target,
987                        entity_set_name_by_logical_name,
988                    )?),
989                    parameters: request.parameters.clone(),
990                }
991            }
992            OrganizationRequest::Delete(request) => {
993                let entity_set_name = entity_set_name_by_logical_name
994                    .get(&request.target.logical_name.to_ascii_lowercase())
995                    .ok_or_else(|| {
996                        format!(
997                            "Entity set metadata not found for '{}'",
998                            request.target.logical_name
999                        )
1000                    })?;
1001
1002                PreparedBatchRequest {
1003                    method: "DELETE",
1004                    path: format!(
1005                        "/api/data/v9.2/{}({})",
1006                        entity_set_name,
1007                        request.target.id.as_hyphenated()
1008                    ),
1009                    body: None,
1010                    parameters: request.parameters.clone(),
1011                }
1012            }
1013        };
1014
1015        Ok(PreparedBatchItem {
1016            prepared_request: prepared,
1017        })
1018    }
1019
1020    fn build_batch_body(&self, boundary: &str, requests: &[PreparedBatchItem]) -> String {
1021        let mut body = String::new();
1022
1023        for (content_id, item) in requests.iter().enumerate() {
1024            // Dataverse expects `$batch` parts in raw HTTP-message form rather than as plain JSON
1025            // fragments, so the multipart body is assembled manually here.
1026            body.push_str(&format!("--{boundary}\r\n"));
1027            body.push_str("Content-Type: application/http\r\n");
1028            body.push_str("Content-Transfer-Encoding: binary\r\n");
1029            body.push_str(&format!("Content-ID: {}\r\n\r\n", content_id + 1));
1030            body.push_str(&format!(
1031                "{} {} HTTP/1.1\r\n",
1032                item.prepared_request.method, item.prepared_request.path
1033            ));
1034            body.push_str("Accept: application/json\r\n");
1035
1036            for (header, value) in item.prepared_request.parameters.headers() {
1037                body.push_str(&format!("{header}: {value}\r\n"));
1038            }
1039
1040            if let Some(payload) = &item.prepared_request.body {
1041                body.push_str("Content-Type: application/json;type=entry\r\n\r\n");
1042                body.push_str(payload);
1043                body.push_str("\r\n");
1044            } else {
1045                body.push_str("\r\n");
1046            }
1047        }
1048
1049        body.push_str(&format!("--{boundary}--\r\n"));
1050        body
1051    }
1052
1053    fn map_batch_response(
1054        &self,
1055        request: &ExecuteMultipleRequest,
1056        parts: Vec<ParsedBatchPart>,
1057    ) -> Result<ExecuteMultipleResponse, String> {
1058        let mut response = ExecuteMultipleResponse::default();
1059
1060        for (part_index, part) in parts.iter().enumerate() {
1061            let Some(source_request) = request.requests.get(part_index) else {
1062                break;
1063            };
1064
1065            if part.status_code >= 400 {
1066                response.responses.push(ExecuteMultipleResponseItem {
1067                    request_index: part_index,
1068                    response: None,
1069                    fault: Some(parse_fault(part)),
1070                });
1071                continue;
1072            }
1073
1074            if request.settings.return_responses {
1075                response.responses.push(ExecuteMultipleResponseItem {
1076                    request_index: part_index,
1077                    response: Some(source_request.success_response(&part.headers)),
1078                    fault: None,
1079                });
1080            }
1081        }
1082
1083        Ok(response)
1084    }
1085}
1086
1087fn ensure_fetch_page_size(fetchxml: &str, page_size: i32) -> Result<String, String> {
1088    if fetch_tag_has_attr(fetchxml, "count")? {
1089        return Ok(fetchxml.to_string());
1090    }
1091
1092    let fetch_start = fetchxml
1093        .find("<fetch")
1094        .ok_or_else(|| "FetchXML must start with a <fetch> element".to_string())?;
1095    let tag_end = fetchxml[fetch_start..]
1096        .find('>')
1097        .ok_or_else(|| "FetchXML <fetch> element is not closed".to_string())?
1098        + fetch_start;
1099
1100    let mut inserted = String::new();
1101    inserted.push_str(&fetchxml[..tag_end]);
1102    inserted.push_str(&format!(" count=\"{page_size}\""));
1103    inserted.push_str(&fetchxml[tag_end..]);
1104    Ok(inserted)
1105}
1106
1107fn normalize_entity_name(value: &str) -> String {
1108    value
1109        .trim_matches(|ch| ch == '[' || ch == ']' || ch == '"' || ch == '`')
1110        .to_ascii_lowercase()
1111}
1112
1113fn parse_uuid_from_uri(value: &str) -> Option<Uuid> {
1114    let start = value.rfind('(')? + 1;
1115    let end = value.rfind(')')?;
1116    Uuid::parse_str(value[start..end].trim_matches('{').trim_matches('}')).ok()
1117}
1118
1119#[cfg(test)]
1120mod tests {
1121    use super::{ensure_fetch_page_size, normalize_entity_name, parse_uuid_from_uri};
1122    use uuid::Uuid;
1123
1124    #[test]
1125    fn ensure_fetch_page_size_adds_count_when_missing() {
1126        let fetchxml = ensure_fetch_page_size("<fetch><entity name=\"account\" /></fetch>", 250)
1127            .expect("should insert count");
1128
1129        assert!(fetchxml.contains("count=\"250\""));
1130    }
1131
1132    #[test]
1133    fn ensure_fetch_page_size_preserves_existing_count() {
1134        let fetchxml =
1135            ensure_fetch_page_size("<fetch count=\"100\"><entity name=\"account\" /></fetch>", 250)
1136                .expect("should preserve count");
1137
1138        assert_eq!(
1139            fetchxml,
1140            "<fetch count=\"100\"><entity name=\"account\" /></fetch>"
1141        );
1142    }
1143
1144    #[test]
1145    fn normalize_entity_name_strips_common_identifier_wrappers() {
1146        assert_eq!(normalize_entity_name("[Account]"), "account");
1147        assert_eq!(normalize_entity_name("\"Contact\""), "contact");
1148        assert_eq!(normalize_entity_name("`Lead`"), "lead");
1149    }
1150
1151    #[test]
1152    fn parse_uuid_from_uri_reads_guid_between_parentheses() {
1153        let parsed = parse_uuid_from_uri(
1154            "https://example.crm.dynamics.com/api/data/v9.2/accounts({aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee})",
1155        );
1156
1157        assert_eq!(
1158            parsed,
1159            Some(Uuid::parse_str("aaaaaaaa-bbbb-cccc-dddd-eeeeeeeeeeee").expect("uuid"))
1160        );
1161    }
1162}