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