1#![allow(dead_code, clippy::type_complexity)]
3
4use std::collections::BTreeMap;
5use std::sync::Arc;
6
7use teaql_core::{
8 DeleteCommand, Entity, EntityDescriptor, InsertCommand, MutationValues, PropertyDescriptor,
9 UpdateCommand, Value,
10};
11
12use crate::entity_status::EntityStatus;
13use crate::{
14 DataServiceError, GraphMutationKind, GraphMutationPlan, GraphNode, GraphOperation,
15 RuntimeError, ScopedCommentNode, TraceScopeToken, sorted_update_fields,
16};
17
18use super::{EntityDataService, helpers::*};
19
20fn recover_trace_or_default(token: &Option<Arc<TraceScopeToken>>) -> Vec<teaql_core::TraceNode> {
21 token
22 .as_ref()
23 .map(|t| t.recover_trace_chain())
24 .unwrap_or_default()
25}
26
27fn resolve_trace_chain(
28 specific: Vec<teaql_core::TraceNode>,
29 fallback: &[teaql_core::TraceNode],
30) -> Vec<teaql_core::TraceNode> {
31 match specific.is_empty() {
32 true => fallback.to_vec(),
33 false => specific,
34 }
35}
36
37impl<'a, E> EntityDataService<'a, E>
38where
39 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
40{
41 pub(crate) async fn save_graph_internal(
42 &self,
43 node: GraphNode,
44 ) -> Result<GraphNode, DataServiceError<E::Error>> {
45 if node.entity != self.entity {
46 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
47 "entity data service {} cannot save graph root {}",
48 self.entity, node.entity
49 ))));
50 }
51 let plan = self.plan_graph(node).await?;
52 self.execute_graph_plan_internal(plan).await
53 }
54
55 pub(crate) async fn save_entity_graph_from_internal(
56 &self,
57 graph: teaql_core::EntityGraph,
58 ) -> Result<GraphNode, DataServiceError<E::Error>> {
59 fn convert(node: teaql_core::EntityGraphNode) -> GraphNode {
60 let mut relations = BTreeMap::new();
61 for (rel_name, child) in node.children {
62 relations
63 .entry(rel_name)
64 .or_insert_with(Vec::new)
65 .push(convert(child));
66 }
67 GraphNode {
68 entity: node.entity_type,
69 values: node.values.into(),
70 relations,
71 operation: match node.operation {
72 teaql_core::EntityGraphOperation::Save => crate::GraphOperation::Upsert,
73 teaql_core::EntityGraphOperation::Delete => crate::GraphOperation::Remove,
74 },
75 comment: node.comment,
76 dirty_fields: None,
77 original_values: None,
78 }
79 }
80 self.save_graph_internal(convert(graph.root)).await
81 }
82
83 pub(crate) async fn save_entity_graph_internal<T>(
84 &self,
85 entity: T,
86 ) -> Result<GraphNode, DataServiceError<E::Error>>
87 where
88 T: Entity,
89 {
90 let node = self
91 .graph_node_from_entity(entity)
92 .map_err(DataServiceError::Runtime)?;
93 self.save_graph_internal(node).await
94 }
95
96 pub(crate) async fn save_entity_internal<T>(
97 &self,
98 entity: T,
99 status: EntityStatus,
100 ) -> Result<GraphNode, DataServiceError<E::Error>>
101 where
102 T: Entity,
103 {
104 if !status.need_persist() {
105 return Ok(GraphNode::new(&self.entity));
106 }
107 if status.is_deleted() {
108 let mut node = self
109 .graph_node_from_entity(entity)
110 .map_err(DataServiceError::Runtime)?;
111 node.operation = GraphOperation::Remove;
112 node.relations.clear();
113 return self.save_graph_internal(node).await;
114 }
115 self.save_entity_graph_internal(entity).await
116 }
117 pub(crate) async fn save_entity_with_comment_internal<T>(
118 &self,
119 entity: T,
120 status: EntityStatus,
121 comment: impl Into<String>,
122 ) -> Result<GraphNode, DataServiceError<E::Error>>
123 where
124 T: Entity,
125 {
126 if status.is_deleted() {
127 let mut node = self
128 .graph_node_from_entity(entity)
129 .map_err(DataServiceError::Runtime)?;
130 node.operation = GraphOperation::Remove;
131 node.relations.clear();
132 node.set_comment(comment);
133 return self.save_graph_internal(node).await;
134 }
135 self.save_entity_graph_with_comment_internal(entity, comment)
136 .await
137 }
138 pub(crate) async fn save_entity_graph_with_comment_internal<T>(
139 &self,
140 entity: T,
141 comment: impl Into<String>,
142 ) -> Result<GraphNode, DataServiceError<E::Error>>
143 where
144 T: Entity,
145 {
146 let mut node = self
147 .graph_node_from_entity(entity)
148 .map_err(DataServiceError::Runtime)?;
149 node.set_comment(comment);
150 self.save_graph_internal(node).await
151 }
152
153 pub(crate) async fn create_entity_graph_with_comment_internal<T>(
157 &self,
158 entity: T,
159 comment: impl Into<String>,
160 ) -> Result<GraphNode, DataServiceError<E::Error>>
161 where
162 T: Entity,
163 {
164 let mut node = self
165 .graph_node_from_entity(entity)
166 .map_err(DataServiceError::Runtime)?;
167 node.operation = GraphOperation::Create;
168 node.set_comment(comment);
169 self.save_graph_internal(node).await
170 }
171
172 pub async fn plan_graph(
173 &self,
174 node: GraphNode,
175 ) -> Result<GraphMutationPlan, DataServiceError<E::Error>> {
176 if node.entity != self.entity {
177 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
178 "entity data service {} cannot plan graph root {}",
179 self.entity, node.entity
180 ))));
181 }
182 let mut node = node;
183 let mut plan = GraphMutationPlan::default();
184 self.collect_graph_plan(&mut node, &mut plan, None, None, false)
185 .await?;
186 plan.planned_root = Some(node);
187 plan.rebuild_batches();
188 Ok(plan)
189 }
190
191 pub(crate) async fn execute_graph_plan_internal(
192 &self,
193 plan: GraphMutationPlan,
194 ) -> Result<GraphNode, DataServiceError<E::Error>> {
195 let Some(mut root) = plan.planned_root else {
196 return Err(DataServiceError::Runtime(RuntimeError::Graph(
197 "graph mutation plan has no planned root".to_owned(),
198 )));
199 };
200
201 for batch in plan.batches {
202 if batch.items.is_empty()
203 || (matches!(batch.kind, GraphMutationKind::Update)
204 && batch.update_fields.is_empty())
205 {
206 continue;
207 }
208 match batch.kind {
209 GraphMutationKind::Create => {
210 let mut cmd = teaql_core::BatchInsertCommand::new(&batch.entity);
211 for item in batch.items {
212 cmd.batch_values.push(item.values);
213 cmd.trace_chains
214 .push(recover_trace_or_default(&item.scope_token));
215 }
216 self.execute_prepared_batch_insert(cmd).await?;
217 }
218 GraphMutationKind::Update => {
219 if batch.update_fields.is_empty() {
220 continue;
221 }
222 let mut cmd =
223 teaql_core::BatchUpdateCommand::new(&batch.entity, batch.update_fields);
224 for item in batch.items {
225 let id = item.values.get("id").cloned().ok_or_else(|| {
226 DataServiceError::Runtime(RuntimeError::Graph(format!(
227 "update item in batch missing id for {}",
228 batch.entity
229 )))
230 })?;
231 let version = item.values.get("version").and_then(|v| match v {
232 teaql_core::Value::I64(n) => Some(*n),
233 _ => None,
234 });
235 cmd.batch_values.push(item.values);
236 cmd.batch_ids.push(id);
237 cmd.batch_expected_versions.push(version);
238 cmd.batch_old_values.push(item.old_values);
239 cmd.trace_chains
240 .push(recover_trace_or_default(&item.scope_token));
241 }
242 self.execute_prepared_batch_update(cmd).await?;
243 }
244 GraphMutationKind::Delete => {
245 for item in batch.items {
247 let id = item.values.get("id").cloned().ok_or_else(|| {
248 DataServiceError::Runtime(RuntimeError::Graph(format!(
249 "delete item in batch missing id for {}",
250 batch.entity
251 )))
252 })?;
253 let mut cmd = teaql_core::DeleteCommand::new(&batch.entity, id);
254 if let Some(teaql_core::Value::I64(version)) = item.values.get("version") {
255 cmd = cmd.expected_version(*version);
256 }
257 let trace_chain = recover_trace_or_default(&item.scope_token);
258 self.delete_scoped_internal(&cmd, trace_chain).await?;
259 }
260 }
261 GraphMutationKind::Reference => {
262 }
264 }
265 }
266
267 if root.operation != GraphOperation::Remove {
268 let descriptor = self
269 .data_service
270 .metadata
271 .context
272 .require_entity(&root.entity)
273 .map_err(DataServiceError::Runtime)?;
274 let id_property = descriptor.id_property().ok_or_else(|| {
275 DataServiceError::Runtime(RuntimeError::Graph(format!(
276 "entity {} has no id property",
277 root.entity
278 )))
279 })?;
280 let id = root.values.get(&id_property.name).cloned().ok_or_else(|| {
281 DataServiceError::Runtime(RuntimeError::Graph(format!(
282 "saved {} missing identity field {}",
283 root.entity, id_property.name
284 )))
285 })?;
286 root.values = self
287 .fetch_graph_current_row_internal(
288 &root.entity,
289 &id_property.name,
290 &id,
291 root.comment
292 .clone()
293 .map(|comment| {
294 vec![teaql_core::TraceNode {
295 kind: teaql_core::TraceKind::AuditReason,
296 entity_type: root.entity.clone(),
297 entity_id: id.try_u64(),
298 comment,
299 }]
300 })
301 .unwrap_or_default(),
302 )
303 .await?
304 .map(Into::into)
305 .ok_or_else(|| {
306 DataServiceError::Runtime(RuntimeError::Graph(format!(
307 "persisted {} record could not be read back",
308 root.entity
309 )))
310 })?;
311 }
312
313 Ok(root)
314 }
315
316 pub fn graph_node_from_entity<T>(&self, entity: T) -> Result<GraphNode, RuntimeError>
317 where
318 T: Entity,
319 {
320 let descriptor = T::entity_descriptor();
321 if descriptor.name != self.entity {
322 return Err(RuntimeError::Graph(format!(
323 "entity data service {} cannot extract graph root {}",
324 self.entity, descriptor.name
325 )));
326 }
327 let dirty_fields = entity.dirty_fields();
330 let original_values = entity.original_values();
331 let is_deleted = entity.is_marked_as_delete();
332 let comment = entity.get_comment();
333 let mut node = self.graph_node_from_values(&descriptor.name, entity.into_values())?;
334 node.dirty_fields = dirty_fields;
335 node.original_values = original_values;
336 if is_deleted {
337 node.operation = GraphOperation::Remove;
338 node.relations.clear();
339 }
340 if let Some(c) = comment {
341 node.set_comment(c);
342 }
343 Ok(node)
344 }
345
346 fn collect_graph_plan<'b, 's: 'b>(
347 &'b self,
348 node: &'b mut GraphNode,
349 plan: &'b mut GraphMutationPlan,
350 parent_scope: Option<&'s ScopedCommentNode<'s>>,
351 parent_token: Option<Arc<TraceScopeToken>>,
352 parent_is_create: bool,
353 ) -> std::pin::Pin<
354 Box<dyn std::future::Future<Output = Result<(), DataServiceError<E::Error>>> + Send + '_>,
355 > {
356 Box::pin(async move {
357 match node.operation {
358 GraphOperation::Reference => {
359 plan.push(
360 node.entity.clone(),
361 GraphMutationKind::Reference,
362 node.values.clone().into(),
363 Vec::new(),
364 parent_token,
365 node.original_values.clone(),
366 );
367 return Ok(());
368 }
369 GraphOperation::Remove => {
370 plan.push(
371 node.entity.clone(),
372 GraphMutationKind::Delete,
373 node.values.clone().into(),
374 Vec::new(),
375 parent_token,
376 node.original_values.clone(),
377 );
378 return Ok(());
379 }
380 GraphOperation::Upsert | GraphOperation::Create => {}
381 }
382
383 let descriptor = self
384 .data_service
385 .metadata
386 .context
387 .require_entity(&node.entity)
388 .map_err(DataServiceError::Runtime)?;
389
390 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
392 parent: parent_scope,
393 track: teaql_core::TraceNode {
394 kind: teaql_core::TraceKind::AuditReason,
395 entity_type: node.entity.clone(),
396 entity_id: node.id().and_then(|v| match v {
397 Value::U64(n) => Some(*n),
398 Value::I64(n) => Some(*n as u64),
399 _ => None,
400 }),
401 comment: c.clone(),
402 },
403 });
404 let active_scope = current_scope.as_ref().or(parent_scope);
405
406 let id_property = descriptor.id_property().cloned();
407 let id = id_property.as_ref().and_then(|property| {
408 node.values
409 .get(&property.name)
410 .filter(|value| !is_unassigned_id_value(value))
411 .cloned()
412 });
413
414 if let Some(id_val) = &id
415 && !plan
416 .visited_nodes
417 .insert((node.entity.clone(), graph_identity_key(id_val)))
418 {
419 return Ok(());
420 }
421
422 let is_create_op = node.operation == GraphOperation::Create
423 || (parent_is_create && node.operation == GraphOperation::Upsert);
424
425 let is_update = match is_create_op {
426 true => false,
427 false => match (id_property.as_ref(), id.as_ref()) {
428 (Some(id_property), Some(id)) => self
429 .fetch_graph_current_row_internal(
430 &node.entity,
431 &id_property.name,
432 id,
433 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
434 )
435 .await?
436 .is_some(),
437 _ => false,
438 },
439 };
440 if !is_update {
441 if let Some(id_property) = id_property.as_ref() {
442 let needs_id = !node.values.contains_key(&id_property.name)
443 || node
444 .values
445 .get(&id_property.name)
446 .is_some_and(is_unassigned_id_value);
447 if needs_id {
448 let id = self
449 .data_service
450 .metadata
451 .context
452 .next_id(&node.entity)
453 .map_err(DataServiceError::Runtime)?;
454 node.values.insert(id_property.name.clone(), Value::U64(id));
455 }
456 }
457 ensure_initial_version(&mut node.values, descriptor);
458 crate::data_service::helpers::ensure_timestamps(&mut node.values, descriptor, true);
459 } else {
460 crate::data_service::helpers::ensure_timestamps(
461 &mut node.values,
462 descriptor,
463 false,
464 );
465 }
466 let update_fields = if is_update {
467 let mut excluded = Vec::new();
468 if let Some(id_property) = id_property.as_ref() {
469 excluded.push(id_property.name.clone());
470 }
471 if let Some(version_property) = descriptor.version_property() {
472 excluded.push(version_property.name.clone());
473 }
474 let mut fields = sorted_update_fields(&node.values, excluded);
475 if let Some(dirty) = &node.dirty_fields {
476 fields.retain(|f| dirty.contains(f));
477 }
478 fields
479 } else {
480 Default::default()
481 };
482
483 let current_token = node
486 .comment
487 .as_ref()
488 .map(|c| {
489 Arc::new(TraceScopeToken {
490 parent: parent_token.clone(),
491 track: teaql_core::TraceNode {
492 kind: teaql_core::TraceKind::AuditReason,
493 entity_type: node.entity.clone(),
494 entity_id: node.id().and_then(|v| match v {
495 Value::U64(n) => Some(*n),
496 Value::I64(n) => Some(*n as u64),
497 _ => None,
498 }),
499 comment: c.clone(),
500 },
501 node_index: plan.next_item_index,
502 })
503 })
504 .or_else(|| parent_token.clone());
505
506 plan.push(
507 node.entity.clone(),
508 GraphMutationKind::for_update(is_update),
509 node.values.clone().into(),
510 update_fields,
511 current_token.clone(),
512 node.original_values.clone(),
513 );
514
515 for (name, children) in &mut node.relations {
516 let relation = descriptor.relation_by_name(name).ok_or_else(|| {
517 DataServiceError::Runtime(RuntimeError::MissingRelation {
518 entity: node.entity.clone(),
519 relation: name.clone(),
520 })
521 })?;
522 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
523 for child in children {
524 ensure_relation_target(&node.entity, name, &relation.target_entity, child)?;
525 child_repo
526 .collect_graph_plan(
527 child,
528 plan,
529 active_scope,
530 current_token.clone(),
531 is_create_op,
532 )
533 .await?;
534 }
535 }
536 Ok(())
537 })
538 }
539
540 fn insert_graph_node_scoped<'b, 's: 'b>(
541 &'b self,
542 mut node: GraphNode,
543 parent_scope: Option<&'s ScopedCommentNode<'s>>,
544 ) -> std::pin::Pin<
545 Box<
546 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
547 + Send
548 + '_,
549 >,
550 > {
551 Box::pin(async move {
552 match node.operation {
553 GraphOperation::Upsert | GraphOperation::Create => {}
554 GraphOperation::Reference => {
555 return self
556 .validate_reference_node(
557 node,
558 parent_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
559 )
560 .await;
561 }
562 GraphOperation::Remove => {
563 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
564 "create graph cannot remove node {}",
565 node.entity
566 ))));
567 }
568 }
569
570 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
572 parent: parent_scope,
573 track: teaql_core::TraceNode {
574 kind: teaql_core::TraceKind::AuditReason,
575 entity_type: node.entity.clone(),
576 entity_id: node.id().and_then(|v| match v {
577 Value::U64(n) => Some(*n),
578 Value::I64(n) => Some(*n as u64),
579 _ => None,
580 }),
581 comment: c.clone(),
582 },
583 });
584 let active_scope = current_scope.as_ref().or(parent_scope);
585
586 let descriptor = self
587 .data_service
588 .metadata
589 .context
590 .require_entity(&node.entity)
591 .map_err(DataServiceError::Runtime)?;
592
593 let mut one_relations = Vec::new();
594 let mut many_relations = Vec::new();
595 for (name, children) in std::mem::take(&mut node.relations) {
596 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
597 DataServiceError::Runtime(RuntimeError::MissingRelation {
598 entity: node.entity.clone(),
599 relation: name.clone(),
600 })
601 })?;
602 match relation.many {
603 true => many_relations.push((name, relation.clone(), children)),
604 false => one_relations.push((name, relation.clone(), children)),
605 }
606 }
607
608 for (name, relation, children) in one_relations {
609 if children.len() > 1 {
610 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
611 "relation {}.{} expects one child, got {}",
612 node.entity,
613 name,
614 children.len()
615 ))));
616 }
617 let mut saved_children = Vec::new();
618 for child in children {
619 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
620 let child_repo = self.scoped_data_service_internal(child.entity.clone());
621 let saved_child = child_repo
622 .insert_graph_node_scoped(child, active_scope)
623 .await?;
624 if relation.attach {
625 let foreign_value = saved_child
626 .values
627 .get(&relation.foreign_key)
628 .cloned()
629 .ok_or_else(|| {
630 DataServiceError::Runtime(RuntimeError::Graph(format!(
631 "saved child {} missing foreign key {} for relation {}.{}",
632 relation.target_entity, relation.foreign_key, node.entity, name
633 )))
634 })?;
635 node.values
636 .insert(relation.local_key.clone(), foreign_value);
637 }
638 saved_children.push(saved_child);
639 }
640 node.relations.insert(name, saved_children);
641 }
642
643 let command = self
644 .prepare_insert_command(&InsertCommand {
645 entity: node.entity.clone(),
646 values: node.values.clone().into(),
647 trace_chain: Vec::new(),
648 })
649 .map_err(DataServiceError::Runtime)?;
650 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
651 self.execute_prepared_insert_with_comment(command.clone(), lineage)
652 .await?;
653 node.values = command.values.into();
654 if let Some(id_property) = descriptor.id_property()
655 && let Some(id) = node.values.get(&id_property.name).cloned()
656 {
657 node.values = self
658 .fetch_graph_current_row_internal(
659 &node.entity,
660 &id_property.name,
661 &id,
662 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
663 )
664 .await?
665 .map(Into::into)
666 .ok_or_else(|| {
667 DataServiceError::Runtime(RuntimeError::Graph(format!(
668 "persisted {} record could not be read back",
669 node.entity
670 )))
671 })?;
672 }
673
674 for (name, relation, children) in many_relations {
675 let local_value =
676 node.values
677 .get(&relation.local_key)
678 .cloned()
679 .ok_or_else(|| {
680 DataServiceError::Runtime(RuntimeError::Graph(format!(
681 "parent {} missing local key {} for relation {}",
682 node.entity, relation.local_key, name
683 )))
684 })?;
685 let mut saved_children = Vec::new();
686 for mut child in children {
687 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
688 if relation.attach {
689 child
690 .values
691 .insert(relation.foreign_key.clone(), local_value.clone());
692 }
693 let child_repo = self.scoped_data_service_internal(child.entity.clone());
694 saved_children.push(
695 child_repo
696 .insert_graph_node_scoped(child, active_scope)
697 .await?,
698 );
699 }
700 node.relations.insert(name, saved_children);
701 }
702
703 Ok(node)
704 })
705 }
706
707 fn upsert_graph_node_scoped<'b, 's: 'b>(
708 &'b self,
709 mut node: GraphNode,
710 parent_scope: Option<&'s ScopedCommentNode<'s>>,
711 ) -> std::pin::Pin<
712 Box<
713 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
714 + Send
715 + '_,
716 >,
717 > {
718 Box::pin(async move {
719 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
721 parent: parent_scope,
722 track: teaql_core::TraceNode {
723 kind: teaql_core::TraceKind::AuditReason,
724 entity_type: node.entity.clone(),
725 entity_id: node.id().and_then(|v| match v {
726 Value::U64(n) => Some(*n),
727 Value::I64(n) => Some(*n as u64),
728 _ => None,
729 }),
730 comment: c.clone(),
731 },
732 });
733 let active_scope = current_scope.as_ref().or(parent_scope);
734
735 match node.operation {
736 GraphOperation::Upsert | GraphOperation::Create => {}
737 GraphOperation::Reference => {
738 return self
739 .validate_reference_node(
740 node,
741 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
742 )
743 .await;
744 }
745 GraphOperation::Remove => {
746 self.validate_remove_node(
747 &node,
748 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
749 )
750 .await?;
751 self.delete_graph_node(&node, parent_scope).await?;
752 return Ok(node);
753 }
754 }
755
756 let descriptor = self
757 .data_service
758 .metadata
759 .context
760 .require_entity(&node.entity)
761 .map_err(DataServiceError::Runtime)?;
762 let Some(id_property) = descriptor.id_property() else {
763 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
764 "entity {} has no id property for graph upsert",
765 node.entity
766 ))));
767 };
768 let Some(id) = node
769 .values
770 .get(&id_property.name)
771 .filter(|value| !is_unassigned_id_value(value))
772 .cloned()
773 else {
774 node.comment = None;
776 return self.insert_graph_node_scoped(node, active_scope).await;
777 };
778
779 if node.operation == GraphOperation::Create
780 || self
781 .fetch_graph_current_row_internal(
782 &node.entity,
783 &id_property.name,
784 &id,
785 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
786 )
787 .await?
788 .is_none()
789 {
790 node.comment = None;
791 return self.insert_graph_node_scoped(node, active_scope).await;
792 }
793
794 let mut one_relations = Vec::new();
795 let mut many_relations = Vec::new();
796 for (name, children) in std::mem::take(&mut node.relations) {
797 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
798 DataServiceError::Runtime(RuntimeError::MissingRelation {
799 entity: node.entity.clone(),
800 relation: name.clone(),
801 })
802 })?;
803 match relation.many {
804 true => many_relations.push((name, relation.clone(), children)),
805 false => one_relations.push((name, relation.clone(), children)),
806 }
807 }
808
809 for (name, relation, children) in one_relations {
810 if children.len() > 1 {
811 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
812 "relation {}.{} expects one child, got {}",
813 node.entity,
814 name,
815 children.len()
816 ))));
817 }
818 let mut saved_children = Vec::new();
819 for child in children {
820 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
821 let child_repo = self.scoped_data_service_internal(child.entity.clone());
822 let saved_child = child_repo
823 .upsert_graph_node_scoped(child, active_scope)
824 .await?;
825 if relation.attach {
826 let foreign_value = saved_child
827 .values
828 .get(&relation.foreign_key)
829 .cloned()
830 .ok_or_else(|| {
831 DataServiceError::Runtime(RuntimeError::Graph(format!(
832 "saved child {} missing foreign key {} for relation {}.{}",
833 relation.target_entity, relation.foreign_key, node.entity, name
834 )))
835 })?;
836 node.values
837 .insert(relation.local_key.clone(), foreign_value);
838 }
839 saved_children.push(saved_child);
840 }
841 node.relations.insert(name, saved_children);
842 }
843
844 let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
845 if !update.values.is_empty() {
846 let prepared_update = self
847 .prepare_update_command(&update)
848 .map_err(DataServiceError::Runtime)?;
849 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
850 self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
851 .await?;
852 for (field, value) in &prepared_update.values {
853 node.values.insert(field.clone(), value.clone());
854 }
855 if let Some(version_property) = descriptor.version_property()
856 && let Some(expected_version) = prepared_update.expected_version
857 {
858 node.values.insert(
859 version_property.name.clone(),
860 Value::I64(expected_version + 1),
861 );
862 }
863 node.values = self
864 .fetch_graph_current_row_internal(
865 &node.entity,
866 &id_property.name,
867 &id,
868 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
869 )
870 .await?
871 .map(Into::into)
872 .ok_or_else(|| {
873 DataServiceError::Runtime(RuntimeError::Graph(format!(
874 "persisted {} record could not be read back",
875 node.entity
876 )))
877 })?;
878 }
879
880 for (name, relation, children) in many_relations {
881 let local_value =
882 node.values
883 .get(&relation.local_key)
884 .cloned()
885 .ok_or_else(|| {
886 DataServiceError::Runtime(RuntimeError::Graph(format!(
887 "parent {} missing local key {} for relation {}",
888 node.entity, relation.local_key, name
889 )))
890 })?;
891 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
892 let child_descriptor = self
893 .data_service
894 .metadata
895 .context
896 .require_entity(&relation.target_entity)
897 .map_err(DataServiceError::Runtime)?;
898 let child_id_property = child_descriptor.id_property().ok_or_else(|| {
899 DataServiceError::Runtime(RuntimeError::Graph(format!(
900 "entity {} has no id property",
901 relation.target_entity
902 )))
903 })?;
904
905 let mut seen = std::collections::BTreeSet::new();
906 let mut saved_children = Vec::new();
907 for mut child in children {
908 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
909 if relation.attach && child.operation != GraphOperation::Reference {
910 child
911 .values
912 .insert(relation.foreign_key.clone(), local_value.clone());
913 }
914 if let Some(child_id) = child
915 .values
916 .get(&child_id_property.name)
917 .filter(|value| !is_unassigned_id_value(value))
918 {
919 let key = graph_identity_key(child_id);
920 if !seen.insert(key.clone()) {
921 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
922 "duplicate child id {key} in relation {}.{}",
923 node.entity, name
924 ))));
925 }
926 }
927 saved_children.push(
928 child_repo
929 .upsert_graph_node_scoped(child, active_scope)
930 .await?,
931 );
932 }
933
934 node.relations.insert(name, saved_children);
935 }
936
937 Ok(node)
938 })
939 }
940
941 async fn validate_reference_node(
942 &self,
943 node: GraphNode,
944 trace_chain: Vec<teaql_core::TraceNode>,
945 ) -> Result<GraphNode, DataServiceError<E::Error>> {
946 if !node.relations.is_empty() {
947 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
948 "reference node {} cannot contain child relations",
949 node.entity
950 ))));
951 }
952 let descriptor = self
953 .data_service
954 .metadata
955 .context
956 .require_entity(&node.entity)
957 .map_err(DataServiceError::Runtime)?;
958 let id_property = descriptor.id_property().ok_or_else(|| {
959 DataServiceError::Runtime(RuntimeError::Graph(format!(
960 "entity {} has no id property for graph reference",
961 node.entity
962 )))
963 })?;
964 let id = node
965 .values
966 .get(&id_property.name)
967 .filter(|value| !is_unassigned_id_value(value))
968 .cloned()
969 .ok_or_else(|| {
970 DataServiceError::Runtime(RuntimeError::Graph(format!(
971 "reference node {} missing id property {}",
972 node.entity, id_property.name
973 )))
974 })?;
975
976 for field in node.values.keys() {
977 if field == &id_property.name {
978 continue;
979 }
980 if descriptor
981 .version_property()
982 .map(|property| field == &property.name)
983 .unwrap_or(false)
984 {
985 continue;
986 }
987 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
988 "reference node {} cannot carry mutable field {}",
989 node.entity, field
990 ))));
991 }
992
993 let current = self
994 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
995 .await?
996 .ok_or_else(|| {
997 DataServiceError::Runtime(RuntimeError::Graph(format!(
998 "reference node {}({}) does not exist",
999 node.entity,
1000 graph_identity_key(&id)
1001 )))
1002 })?;
1003
1004 if let Some(version_property) = descriptor.version_property()
1005 && let Some(Value::I64(existing_version)) = current.get(&version_property.name)
1006 {
1007 if *existing_version < 0 {
1008 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1009 "reference node {}({}) is deleted",
1010 node.entity,
1011 graph_identity_key(&id)
1012 ))));
1013 }
1014 if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
1015 && expected_version != existing_version
1016 {
1017 println!(
1018 "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
1019 node.entity, expected_version, existing_version
1020 );
1021 return Err(DataServiceError::Runtime(
1022 RuntimeError::OptimisticLockConflict {
1023 entity: node.entity,
1024 id: graph_identity_key(&id),
1025 },
1026 ));
1027 }
1028 }
1029
1030 Ok(GraphNode {
1031 entity: node.entity,
1032 values: current.into(),
1033 relations: BTreeMap::new(),
1034 operation: GraphOperation::Reference,
1035 comment: None,
1036 dirty_fields: None,
1037 original_values: None,
1038 })
1039 }
1040
1041 async fn validate_remove_node(
1042 &self,
1043 node: &GraphNode,
1044 trace_chain: Vec<teaql_core::TraceNode>,
1045 ) -> Result<(), DataServiceError<E::Error>> {
1046 if !node.relations.is_empty() {
1047 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1048 "remove node {} cannot contain child relations",
1049 node.entity
1050 ))));
1051 }
1052 let descriptor = self
1053 .data_service
1054 .metadata
1055 .context
1056 .require_entity(&node.entity)
1057 .map_err(DataServiceError::Runtime)?;
1058 let id_property = descriptor.id_property().ok_or_else(|| {
1059 DataServiceError::Runtime(RuntimeError::Graph(format!(
1060 "entity {} has no id property for graph remove",
1061 node.entity
1062 )))
1063 })?;
1064 let id = node
1065 .values
1066 .get(&id_property.name)
1067 .filter(|value| !is_unassigned_id_value(value))
1068 .cloned()
1069 .ok_or_else(|| {
1070 DataServiceError::Runtime(RuntimeError::Graph(format!(
1071 "remove node {} missing id property {}",
1072 node.entity, id_property.name
1073 )))
1074 })?;
1075 let current = self
1076 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
1077 .await?
1078 .ok_or_else(|| {
1079 DataServiceError::Runtime(RuntimeError::Graph(format!(
1080 "remove node {}({}) does not exist",
1081 node.entity,
1082 graph_identity_key(&id)
1083 )))
1084 })?;
1085 if let Some(version_property) = descriptor.version_property()
1086 && let Some(Value::I64(existing_version)) = current.get(&version_property.name)
1087 && *existing_version < 0
1088 {
1089 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1090 "remove node {}({}) is already deleted",
1091 node.entity,
1092 graph_identity_key(&id)
1093 ))));
1094 }
1095 Ok(())
1096 }
1097
1098 fn graph_node_from_values(
1099 &self,
1100 entity: &str,
1101 values: MutationValues,
1102 ) -> Result<GraphNode, RuntimeError> {
1103 let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1104 let mut node = GraphNode::new(entity);
1105
1106 for (field, value) in values {
1107 if field == "_comment" {
1108 if let Value::Text(comment) = value {
1109 node.set_comment(comment);
1110 }
1111 continue;
1112 }
1113 if field == "_dirty_fields" {
1114 if let Value::List(fields) = value {
1115 let mut dirty = std::collections::BTreeSet::new();
1116 for f in fields {
1117 if let Value::Text(t) = f {
1118 dirty.insert(t);
1119 }
1120 }
1121 node.dirty_fields = Some(dirty);
1122 }
1123 continue;
1124 }
1125 if field == "_original_values" {
1126 if let Value::Object(orig) = value {
1127 node.original_values = Some(orig.into());
1128 }
1129 continue;
1130 }
1131 if field == "_is_new" {
1132 if matches!(value, Value::Bool(true)) {
1133 node.operation = GraphOperation::Create;
1134 }
1135 continue;
1136 }
1137 if field == "_is_deleted" {
1138 if matches!(value, Value::Bool(true)) {
1139 node.operation = GraphOperation::Remove;
1140 }
1141 continue;
1142 }
1143 let Some(relation) = descriptor.relation_by_name(&field) else {
1144 node.values.insert(field, value);
1145 continue;
1146 };
1147
1148 match value {
1149 Value::Null => {
1150 node.relations.entry(field).or_default();
1151 }
1152 Value::Object(record) => {
1153 let child =
1154 self.graph_node_from_values(&relation.target_entity, record.into())?;
1155 node.relations.entry(field).or_default().push(child);
1156 }
1157 Value::List(values) => {
1158 let children = node.relations.entry(field.clone()).or_default();
1159 for value in values {
1160 let Value::Object(record) = value else {
1161 return Err(RuntimeError::Graph(format!(
1162 "relation {}.{} expects object children, got {:?}",
1163 entity, field, value
1164 )));
1165 };
1166 children.push(
1167 self.graph_node_from_values(&relation.target_entity, record.into())?,
1168 );
1169 }
1170 }
1171 other => {
1172 return Err(RuntimeError::Graph(format!(
1173 "relation {}.{} expects object/list/null, got {:?}",
1174 entity, field, other
1175 )));
1176 }
1177 }
1178 }
1179
1180 Ok(node)
1181 }
1182
1183 fn graph_update_command(
1184 &self,
1185 node: &mut GraphNode,
1186 descriptor: &EntityDescriptor,
1187 id_property: &PropertyDescriptor,
1188 id: &Value,
1189 ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1190 crate::mark_entity_status(&mut node.values, crate::CheckObjectStatus::Update);
1191 let check_result = self
1192 .data_service
1193 .metadata
1194 .context
1195 .check_and_fix_values(&node.entity, &mut node.values);
1196 crate::clear_entity_status(&mut node.values);
1197 check_result.map_err(DataServiceError::Runtime)?;
1198
1199 let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1200 command.old_values = node.original_values.clone();
1201 if let Some(version_property) = descriptor.version_property()
1202 && let Some(Value::I64(version)) = node.values.get(&version_property.name)
1203 {
1204 command = command.expected_version(*version);
1205 }
1206 for property in descriptor.properties.iter().filter(|property| {
1210 !property.is_id
1211 && !property.is_version
1212 && property.name != id_property.name
1213 && match &node.dirty_fields {
1214 Some(dirty) => dirty.contains(&property.name),
1215 None => node.values.contains_key(&property.name),
1216 }
1217 }) {
1218 if let Some(value) = node.values.get(&property.name) {
1219 command.values.insert(property.name.clone(), value.clone());
1220 }
1221 }
1222 Ok(command)
1223 }
1224
1225 fn delete_graph_node<'b, 's: 'b>(
1226 &'b self,
1227 node: &'b GraphNode,
1228 parent_scope: Option<&'s ScopedCommentNode<'s>>,
1229 ) -> std::pin::Pin<
1230 Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1231 > {
1232 Box::pin(async move {
1233 let descriptor = self
1234 .data_service
1235 .metadata
1236 .context
1237 .require_entity(&node.entity)
1238 .map_err(DataServiceError::Runtime)?;
1239 let id_property = descriptor.id_property().ok_or_else(|| {
1240 DataServiceError::Runtime(RuntimeError::Graph(format!(
1241 "entity {} has no id property for graph remove",
1242 node.entity
1243 )))
1244 })?;
1245 let id = node
1246 .values
1247 .get(&id_property.name)
1248 .filter(|value| !is_unassigned_id_value(value))
1249 .cloned()
1250 .ok_or_else(|| {
1251 DataServiceError::Runtime(RuntimeError::Graph(format!(
1252 "remove node {} missing id property {}",
1253 node.entity, id_property.name
1254 )))
1255 })?;
1256 let mut delete = DeleteCommand::new(node.entity.clone(), id);
1257 if let Some(version_property) = descriptor.version_property()
1258 && let Some(Value::I64(version)) = node.values.get(&version_property.name)
1259 {
1260 delete = delete.expected_version(*version);
1261 }
1262
1263 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1265 parent: parent_scope,
1266 track: teaql_core::TraceNode {
1267 kind: teaql_core::TraceKind::AuditReason,
1268 entity_type: node.entity.clone(),
1269 entity_id: node.id().and_then(|v| match v {
1270 Value::U64(n) => Some(*n),
1271 Value::I64(n) => Some(*n as u64),
1272 _ => None,
1273 }),
1274 comment: c.clone(),
1275 },
1276 });
1277 let active_scope = current_scope.as_ref().or(parent_scope);
1278 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1279
1280 self.delete_scoped_internal(&delete, lineage).await
1281 })
1282 }
1283
1284 pub(crate) async fn fetch_graph_current_row_internal(
1285 &self,
1286 entity: &str,
1287 id_property: &str,
1288 id: &teaql_core::Value,
1289 trace_chain: Vec<teaql_core::TraceNode>,
1290 ) -> Result<Option<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1291 let mut query = teaql_core::SelectQuery::new(entity)
1292 .filter(teaql_core::Expr::eq(id_property, id.clone()));
1293 query.trace_chain = trace_chain;
1294 let mut rows = self
1295 .scoped_data_service_internal(entity.to_owned())
1296 .fetch_all_internal(&query)
1297 .await?;
1298 Ok(rows.pop())
1299 }
1300
1301 pub(crate) fn order_new_ledger_keys(
1302 &self,
1303 keys: impl IntoIterator<Item = crate::EntityKey>,
1304 changes: &std::collections::BTreeMap<crate::EntityKey, crate::EntityValues>,
1305 ) -> Result<Vec<crate::EntityKey>, RuntimeError> {
1306 let keys = keys.into_iter().collect::<std::collections::BTreeSet<_>>();
1307 let mut incoming = keys
1308 .iter()
1309 .cloned()
1310 .map(|key| (key, 0_usize))
1311 .collect::<std::collections::BTreeMap<_, _>>();
1312 let mut outgoing = std::collections::BTreeMap::<
1313 crate::EntityKey,
1314 std::collections::BTreeSet<crate::EntityKey>,
1315 >::new();
1316
1317 let add_dependency =
1318 |parent: &crate::EntityKey,
1319 child: &crate::EntityKey,
1320 incoming: &mut std::collections::BTreeMap<crate::EntityKey, usize>,
1321 outgoing: &mut std::collections::BTreeMap<
1322 crate::EntityKey,
1323 std::collections::BTreeSet<crate::EntityKey>,
1324 >| {
1325 if parent != child
1326 && outgoing
1327 .entry(parent.clone())
1328 .or_default()
1329 .insert(child.clone())
1330 {
1331 *incoming
1332 .get_mut(child)
1333 .expect("every inserted key has an indegree") += 1;
1334 }
1335 };
1336
1337 for source in &keys {
1338 let descriptor = self
1339 .data_service
1340 .metadata
1341 .context
1342 .require_entity(source.entity.as_ref())?;
1343 let Some(values) = changes.get(source) else {
1344 continue;
1345 };
1346 for relation in &descriptor.relations {
1347 let Some(source_value) = values
1348 .get(&relation.local_key)
1349 .filter(|value| !matches!(value, Value::Null | Value::TypedNull(_)))
1350 else {
1351 continue;
1352 };
1353
1354 for target in keys
1355 .iter()
1356 .filter(|target| target.entity.as_ref() == relation.target_entity)
1357 {
1358 let target_value = changes
1359 .get(target)
1360 .and_then(|record| record.get(&relation.foreign_key))
1361 .or_else(|| (relation.foreign_key == "id").then_some(&target.id));
1362 if target_value != Some(source_value) {
1363 continue;
1364 }
1365 if relation.many {
1366 add_dependency(source, target, &mut incoming, &mut outgoing);
1367 } else {
1368 add_dependency(target, source, &mut incoming, &mut outgoing);
1369 }
1370 }
1371 }
1372 }
1373
1374 let mut ready = incoming
1375 .iter()
1376 .filter(|(_, count)| **count == 0)
1377 .map(|(key, _)| key.clone())
1378 .collect::<std::collections::BTreeSet<_>>();
1379 let mut ordered = Vec::with_capacity(keys.len());
1380 while let Some(key) = ready.pop_first() {
1381 ordered.push(key.clone());
1382 if let Some(dependents) = outgoing.get(&key) {
1383 for dependent in dependents {
1384 let count = incoming
1385 .get_mut(dependent)
1386 .expect("every dependent key has an indegree");
1387 *count -= 1;
1388 if *count == 0 {
1389 ready.insert(dependent.clone());
1390 }
1391 }
1392 }
1393 }
1394
1395 if ordered.len() != keys.len() {
1396 let cyclic_entities = incoming
1397 .into_iter()
1398 .filter(|(_, count)| *count > 0)
1399 .map(|(key, _)| key.entity.into_owned())
1400 .collect::<std::collections::BTreeSet<_>>()
1401 .into_iter()
1402 .collect::<Vec<_>>()
1403 .join(", ");
1404 return Err(RuntimeError::Graph(format!(
1405 "new entity graph contains a required foreign-key cycle among: {cyclic_entities}"
1406 )));
1407 }
1408 Ok(ordered)
1409 }
1410
1411 pub(crate) async fn execute_ledger_plan_internal(
1412 &self,
1413 root: crate::EntityRuntimeState,
1414 locations: &std::collections::BTreeMap<crate::EntityKey, crate::ObjectLocation>,
1415 ) -> Result<std::collections::BTreeMap<crate::EntityKey, Value>, DataServiceError<E::Error>>
1416 {
1417 if let Some(error) = root.first_composition_error() {
1418 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1419 "generated entity graph attachment failed before save: {error}"
1420 ))));
1421 }
1422 let mut generated_ids = std::collections::BTreeMap::new();
1423 let comment = root.get_comment();
1424 let trace_chain = comment
1425 .map(|c| {
1426 vec![teaql_core::TraceNode {
1427 kind: teaql_core::TraceKind::AuditReason,
1428 entity_type: self.entity.clone(),
1429 entity_id: None,
1430 comment: c,
1431 }]
1432 })
1433 .unwrap_or_default();
1434
1435 let deleted_keys = root.deleted_keys();
1436 let new_keys = root.new_keys();
1437 let change_set = root.current_change_set();
1438
1439 for key in &deleted_keys {
1443 if new_keys.contains(key) {
1444 continue;
1445 }
1446 let descriptor = self
1447 .data_service
1448 .metadata
1449 .context
1450 .require_entity(&key.entity)
1451 .map_err(DataServiceError::Runtime)?;
1452 if descriptor.version_property().is_some() && root.get_original_version(key).is_none() {
1453 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1454 "cannot delete {}({:?}) without its loaded original version; load the full entity before mutation",
1455 key.entity, key.id
1456 ))));
1457 }
1458 }
1459
1460 let mut checked_changes = std::collections::BTreeMap::new();
1467 for (key, record) in change_set.changes() {
1468 if deleted_keys.contains(key) {
1469 continue;
1470 }
1471 let mut checked: crate::EntityValues = record.clone().into();
1472 checked
1473 .entry("id".to_owned())
1474 .or_insert_with(|| key.id.clone());
1475 checked_changes.insert(key.clone(), checked);
1476 }
1477
1478 let mut update_batches: std::collections::BTreeMap<
1482 (String, String),
1483 Vec<crate::EntityKey>,
1484 > = std::collections::BTreeMap::new();
1485 let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1486 std::collections::BTreeMap::new();
1487
1488 for (key, record) in &checked_changes {
1489 if deleted_keys.contains(key) {
1490 continue;
1491 }
1492 let mut is_new = new_keys.contains(key);
1493
1494 if !is_new {
1495 let descriptor = self
1496 .data_service
1497 .metadata
1498 .context
1499 .require_entity(&key.entity)
1500 .map_err(DataServiceError::Runtime)?;
1501 let id_property = descriptor.id_property().ok_or_else(|| {
1502 DataServiceError::Runtime(RuntimeError::Graph(format!(
1503 "entity {} has no id property",
1504 key.entity
1505 )))
1506 })?;
1507 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1508 let current_row = self
1509 .fetch_graph_current_row_internal(
1510 &key.entity,
1511 &id_property.name,
1512 &key.id,
1513 my_trace,
1514 )
1515 .await?;
1516 if current_row.is_none() {
1517 is_new = true;
1518 } else if descriptor.version_property().is_some()
1519 && root.get_original_version(key).is_none()
1520 {
1521 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1522 "cannot update {}({:?}) without its loaded original version; load the full entity before mutation",
1523 key.entity, key.id
1524 ))));
1525 }
1526 }
1527
1528 match is_new {
1529 true => {
1530 insert_batches
1531 .entry(key.entity.to_string())
1532 .or_default()
1533 .push(key.clone());
1534 }
1535 false => {
1536 let mut fields: Vec<String> = record.keys().cloned().collect();
1537 fields.sort();
1538 let signature = fields.join(",");
1539 update_batches
1540 .entry((key.entity.to_string(), signature))
1541 .or_default()
1542 .push(key.clone());
1543 }
1544 }
1545 }
1546
1547 let ordered_insert_keys = self
1548 .order_new_ledger_keys(insert_batches.values().flatten().cloned(), &checked_changes)
1549 .map_err(DataServiceError::Runtime)?;
1550 let mut ordered_insert_batches = Vec::<(String, Vec<crate::EntityKey>)>::new();
1551 for key in ordered_insert_keys {
1552 let entity = key.entity.to_string();
1553 if let Some((batch_entity, keys)) = ordered_insert_batches.last_mut()
1554 && *batch_entity == entity
1555 {
1556 keys.push(key);
1557 } else {
1558 ordered_insert_batches.push((entity, vec![key]));
1559 }
1560 }
1561
1562 let mut prepared_insert_batches = Vec::new();
1563 for (entity, keys) in ordered_insert_batches {
1564 let descriptor = self
1565 .data_service
1566 .metadata
1567 .context
1568 .require_entity(&entity)
1569 .map_err(DataServiceError::Runtime)?;
1570 let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.name);
1571 let mut traces = Vec::new();
1572 for key in &keys {
1573 let record = checked_changes.get(key).unwrap();
1574 let mut db_record = crate::EntityValues::new();
1575 let mut real_id = key.id.clone();
1576 if crate::data_service::helpers::is_unassigned_id_value(&real_id) {
1577 let gen_id = self
1578 .data_service
1579 .metadata
1580 .context
1581 .next_id(&entity)
1582 .map_err(DataServiceError::Runtime)?;
1583 real_id = Value::U64(gen_id);
1584 generated_ids.insert(key.clone(), real_id.clone());
1585 }
1586 db_record.insert("id".to_owned(), real_id);
1587 for (field, value) in record {
1588 if field == "id" {
1589 continue;
1590 }
1591 db_record.insert(field.clone(), value.clone());
1592 }
1593 crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1594 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, true);
1595 let location = locations.get(key).cloned().unwrap_or_default();
1596 self.data_service
1597 .metadata
1598 .context
1599 .validate_required_create_payload(&entity, &db_record, &location)
1600 .map_err(DataServiceError::Runtime)?;
1601 cmd.batch_values.push(db_record.into());
1602 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1603 traces.push(my_trace);
1604 }
1605 cmd.trace_chains = traces;
1606 prepared_insert_batches.push(cmd);
1607 }
1608
1609 for key in &deleted_keys {
1613 if new_keys.contains(key) {
1614 continue;
1615 }
1616 let id = key.id.clone();
1617 let mut cmd = teaql_core::DeleteCommand::new(key.entity.as_ref(), id);
1618 if let Some(version) = root.get_original_version(key) {
1619 cmd = cmd.expected_version(version);
1620 }
1621 cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1622 self.delete_internal(&cmd).await?;
1623 }
1624
1625 for cmd in prepared_insert_batches {
1626 self.execute_prepared_batch_insert(cmd).await?;
1627 }
1628
1629 let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1630 update_order.sort();
1631
1632 for signature in update_order {
1633 let keys = update_batches.get(&signature).unwrap();
1634 let descriptor = self
1635 .data_service
1636 .metadata
1637 .context
1638 .require_entity(&signature.0)
1639 .map_err(DataServiceError::Runtime)?;
1640 let mut update_fields: Vec<String> =
1641 signature.1.split(',').map(|s| s.to_string()).collect();
1642 if descriptor
1643 .properties
1644 .iter()
1645 .any(|p| p.name == "update_time")
1646 && !update_fields.contains(&"update_time".to_owned())
1647 {
1648 update_fields.push("update_time".to_owned());
1649 }
1650 if let Some(version_property) = descriptor.version_property()
1651 && !update_fields.contains(&version_property.name)
1652 {
1653 update_fields.push(version_property.name.clone());
1654 }
1655 let mut cmd = teaql_core::BatchUpdateCommand::new(&descriptor.name, update_fields);
1656 let mut traces = Vec::new();
1657 for key in keys {
1658 let record = checked_changes.get(key).unwrap();
1659 let mut db_record = crate::EntityValues::new();
1660 db_record.insert("id".to_owned(), key.id.clone());
1661 for (field, value) in record {
1662 if field == "id" {
1663 continue;
1664 }
1665 db_record.insert(field.clone(), value.clone());
1666 }
1667 crate::data_service::helpers::increment_version(
1668 &mut db_record,
1669 descriptor,
1670 root.get_original_version(key),
1671 );
1672 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, false);
1673 cmd.batch_values.push(db_record.into());
1674 cmd.batch_ids.push(key.id.clone());
1675 cmd.batch_expected_versions
1676 .push(root.get_original_version(key));
1677 cmd.batch_old_values.push(None); let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1679 traces.push(my_trace);
1680 }
1681 cmd.trace_chains = traces;
1682 self.execute_prepared_batch_update(cmd).await?;
1683 }
1684
1685 Ok(generated_ids)
1686 }
1687}