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