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#[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
85pub 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_cache: Mutex<Option<Vec<EntityDefinition>>>,
95 entity_attributes_cache: Mutex<HashMap<String, Vec<EntityAttribute>>>,
98 log_level: LogLevel,
99}
100
101impl ServiceClient {
102 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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}