1use std::collections::BTreeMap;
2use std::sync::Arc;
3
4use teaql_core::{
5 DeleteCommand, Entity, EntityDescriptor, Expr, InsertCommand, PropertyDescriptor, Record,
6 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.record,
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);
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);
237 cmd.batch_ids.push(id);
238 cmd.batch_expected_versions.push(version);
239 cmd.batch_old_values.push(item.old_values);
240 cmd.trace_chains
241 .push(recover_trace_or_default(&item.scope_token));
242 }
243 self.execute_prepared_batch_update(cmd).await?;
244 }
245 GraphMutationKind::Delete => {
246 for item in batch.items {
248 let id = item.values.get("id").cloned().ok_or_else(|| {
249 DataServiceError::Runtime(RuntimeError::Graph(format!(
250 "delete item in batch missing id for {}",
251 batch.entity
252 )))
253 })?;
254 let mut cmd = teaql_core::DeleteCommand::new(&batch.entity, id);
255 if let Some(teaql_core::Value::I64(version)) = item.values.get("version") {
256 cmd = cmd.expected_version(*version);
257 }
258 let trace_chain = recover_trace_or_default(&item.scope_token);
259 self.delete_scoped_internal(&cmd, trace_chain).await?;
260 }
261 }
262 GraphMutationKind::Reference => {
263 }
265 }
266 }
267
268 Ok(root)
269 }
270
271 pub fn graph_node_from_entity<T>(&self, entity: T) -> Result<GraphNode, RuntimeError>
272 where
273 T: Entity,
274 {
275 let descriptor = T::entity_descriptor();
276 if descriptor.name != self.entity {
277 return Err(RuntimeError::Graph(format!(
278 "entity data service {} cannot extract graph root {}",
279 self.entity, descriptor.name
280 )));
281 }
282 let dirty_fields = entity.dirty_fields();
285 let original_values = entity.original_values();
286 let is_deleted = entity.is_marked_as_delete();
287 let comment = entity.get_comment();
288 let mut node = self.graph_node_from_record(&descriptor.name, entity.into_record())?;
289 node.dirty_fields = dirty_fields;
290 node.original_values = original_values;
291 if is_deleted {
292 node.operation = GraphOperation::Remove;
293 node.relations.clear();
294 }
295 if let Some(c) = comment {
296 node.set_comment(c);
297 }
298 Ok(node)
299 }
300
301 fn collect_graph_plan<'b, 's: 'b>(
302 &'b self,
303 node: &'b mut GraphNode,
304 plan: &'b mut GraphMutationPlan,
305 parent_scope: Option<&'s ScopedCommentNode<'s>>,
306 parent_token: Option<Arc<TraceScopeToken>>,
307 parent_is_create: bool,
308 ) -> std::pin::Pin<
309 Box<dyn std::future::Future<Output = Result<(), DataServiceError<E::Error>>> + Send + '_>,
310 > {
311 Box::pin(async move {
312 match node.operation {
313 GraphOperation::Reference => {
314 plan.push(
315 node.entity.clone(),
316 GraphMutationKind::Reference,
317 node.values.clone(),
318 Vec::new(),
319 parent_token,
320 node.original_values.clone(),
321 );
322 return Ok(());
323 }
324 GraphOperation::Remove => {
325 plan.push(
326 node.entity.clone(),
327 GraphMutationKind::Delete,
328 node.values.clone(),
329 Vec::new(),
330 parent_token,
331 node.original_values.clone(),
332 );
333 return Ok(());
334 }
335 GraphOperation::Upsert | GraphOperation::Create => {}
336 }
337
338 let descriptor = self
339 .data_service
340 .metadata
341 .context
342 .require_entity(&node.entity)
343 .map_err(DataServiceError::Runtime)?;
344
345 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
347 parent: parent_scope,
348 track: teaql_core::TraceNode {
349 entity_type: node.entity.clone(),
350 entity_id: node.id().and_then(|v| match v {
351 Value::U64(n) => Some(*n),
352 Value::I64(n) => Some(*n as u64),
353 _ => None,
354 }),
355 comment: c.clone(),
356 },
357 });
358 let active_scope = current_scope.as_ref().or(parent_scope);
359
360 let id_property = descriptor.id_property().cloned();
361 let id = id_property.as_ref().and_then(|property| {
362 node.values
363 .get(&property.name)
364 .filter(|value| !is_unassigned_id_value(value))
365 .cloned()
366 });
367
368 if let Some(id_val) = &id {
369 if !plan
370 .visited_nodes
371 .insert((node.entity.clone(), graph_identity_key(id_val)))
372 {
373 return Ok(());
374 }
375 }
376
377 let is_create_op = node.operation == GraphOperation::Create
378 || (parent_is_create && node.operation == GraphOperation::Upsert);
379
380 let is_update = match is_create_op {
381 true => false,
382 false => match (id_property.as_ref(), id.as_ref()) {
383 (Some(id_property), Some(id)) => self
384 .fetch_graph_current_row_internal(
385 &node.entity,
386 &id_property.name,
387 id,
388 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
389 )
390 .await?
391 .is_some(),
392 _ => false,
393 },
394 };
395 if !is_update {
396 if let Some(id_property) = id_property.as_ref() {
397 let needs_id = !node.values.contains_key(&id_property.name)
398 || node
399 .values
400 .get(&id_property.name)
401 .is_some_and(is_unassigned_id_value);
402 if needs_id {
403 let id = self
404 .data_service
405 .metadata
406 .context
407 .next_id(&node.entity)
408 .map_err(DataServiceError::Runtime)?;
409 node.values.insert(id_property.name.clone(), Value::U64(id));
410 }
411 }
412 ensure_initial_version(&mut node.values, descriptor);
413 }
414 let update_fields = is_update
415 .then(|| {
416 let mut excluded = Vec::new();
417 if let Some(id_property) = id_property.as_ref() {
418 excluded.push(id_property.name.clone());
419 }
420 if let Some(version_property) = descriptor.version_property() {
421 excluded.push(version_property.name.clone());
422 }
423 let mut fields = sorted_update_fields(&node.values, excluded);
424 if let Some(dirty) = &node.dirty_fields {
425 fields.retain(|f| dirty.contains(f));
426 }
427 fields
428 })
429 .unwrap_or_default();
430
431 let current_token = node
434 .comment
435 .as_ref()
436 .map(|c| {
437 Arc::new(TraceScopeToken {
438 parent: parent_token.clone(),
439 track: teaql_core::TraceNode {
440 entity_type: node.entity.clone(),
441 entity_id: node.id().and_then(|v| match v {
442 Value::U64(n) => Some(*n),
443 Value::I64(n) => Some(*n as u64),
444 _ => None,
445 }),
446 comment: c.clone(),
447 },
448 node_index: plan.next_item_index,
449 })
450 })
451 .or_else(|| parent_token.clone());
452
453 plan.push(
454 node.entity.clone(),
455 GraphMutationKind::for_update(is_update),
456 node.values.clone(),
457 update_fields,
458 current_token.clone(),
459 node.original_values.clone(),
460 );
461
462 for (name, children) in &mut node.relations {
463 let relation = descriptor.relation_by_name(name).ok_or_else(|| {
464 DataServiceError::Runtime(RuntimeError::MissingRelation {
465 entity: node.entity.clone(),
466 relation: name.clone(),
467 })
468 })?;
469 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
470 for child in children {
471 ensure_relation_target(&node.entity, name, &relation.target_entity, child)?;
472 child_repo
473 .collect_graph_plan(
474 child,
475 plan,
476 active_scope,
477 current_token.clone(),
478 is_create_op,
479 )
480 .await?;
481 }
482 }
483 Ok(())
484 })
485 }
486
487 fn insert_graph_node_scoped<'b, 's: 'b>(
488 &'b self,
489 mut node: GraphNode,
490 parent_scope: Option<&'s ScopedCommentNode<'s>>,
491 ) -> std::pin::Pin<
492 Box<
493 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
494 + Send
495 + '_,
496 >,
497 > {
498 Box::pin(async move {
499 match node.operation {
500 GraphOperation::Upsert | GraphOperation::Create => {}
501 GraphOperation::Reference => {
502 return self
503 .validate_reference_node(
504 node,
505 parent_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
506 )
507 .await;
508 }
509 GraphOperation::Remove => {
510 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
511 "create graph cannot remove node {}",
512 node.entity
513 ))));
514 }
515 }
516
517 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
519 parent: parent_scope,
520 track: teaql_core::TraceNode {
521 entity_type: node.entity.clone(),
522 entity_id: node.id().and_then(|v| match v {
523 Value::U64(n) => Some(*n),
524 Value::I64(n) => Some(*n as u64),
525 _ => None,
526 }),
527 comment: c.clone(),
528 },
529 });
530 let active_scope = current_scope.as_ref().or(parent_scope);
531
532 let descriptor = self
533 .data_service
534 .metadata
535 .context
536 .require_entity(&node.entity)
537 .map_err(DataServiceError::Runtime)?;
538
539 let mut one_relations = Vec::new();
540 let mut many_relations = Vec::new();
541 for (name, children) in std::mem::take(&mut node.relations) {
542 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
543 DataServiceError::Runtime(RuntimeError::MissingRelation {
544 entity: node.entity.clone(),
545 relation: name.clone(),
546 })
547 })?;
548 match relation.many {
549 true => many_relations.push((name, relation.clone(), children)),
550 false => one_relations.push((name, relation.clone(), children)),
551 }
552 }
553
554 for (name, relation, children) in one_relations {
555 if children.len() > 1 {
556 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
557 "relation {}.{} expects one child, got {}",
558 node.entity,
559 name,
560 children.len()
561 ))));
562 }
563 let mut saved_children = Vec::new();
564 for child in children {
565 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
566 let child_repo = self.scoped_data_service_internal(child.entity.clone());
567 let saved_child = child_repo
568 .insert_graph_node_scoped(child, active_scope)
569 .await?;
570 if relation.attach {
571 let foreign_value = saved_child
572 .values
573 .get(&relation.foreign_key)
574 .cloned()
575 .ok_or_else(|| {
576 DataServiceError::Runtime(RuntimeError::Graph(format!(
577 "saved child {} missing foreign key {} for relation {}.{}",
578 relation.target_entity, relation.foreign_key, node.entity, name
579 )))
580 })?;
581 node.values
582 .insert(relation.local_key.clone(), foreign_value);
583 }
584 saved_children.push(saved_child);
585 }
586 node.relations.insert(name, saved_children);
587 }
588
589 let command = self
590 .prepare_insert_command(&InsertCommand {
591 entity: node.entity.clone(),
592 values: node.values.clone(),
593 trace_chain: Vec::new(),
594 })
595 .map_err(DataServiceError::Runtime)?;
596 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
597 self.execute_prepared_insert_with_comment(command.clone(), lineage)
598 .await?;
599 node.values = command.values;
600
601 for (name, relation, children) in many_relations {
602 let local_value =
603 node.values
604 .get(&relation.local_key)
605 .cloned()
606 .ok_or_else(|| {
607 DataServiceError::Runtime(RuntimeError::Graph(format!(
608 "parent {} missing local key {} for relation {}",
609 node.entity, relation.local_key, name
610 )))
611 })?;
612 let mut saved_children = Vec::new();
613 for mut child in children {
614 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
615 if relation.attach {
616 child
617 .values
618 .insert(relation.foreign_key.clone(), local_value.clone());
619 }
620 let child_repo = self.scoped_data_service_internal(child.entity.clone());
621 saved_children.push(
622 child_repo
623 .insert_graph_node_scoped(child, active_scope)
624 .await?,
625 );
626 }
627 node.relations.insert(name, saved_children);
628 }
629
630 Ok(node)
631 })
632 }
633
634 fn upsert_graph_node_scoped<'b, 's: 'b>(
635 &'b self,
636 mut node: GraphNode,
637 parent_scope: Option<&'s ScopedCommentNode<'s>>,
638 ) -> std::pin::Pin<
639 Box<
640 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
641 + Send
642 + '_,
643 >,
644 > {
645 Box::pin(async move {
646 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
648 parent: parent_scope,
649 track: teaql_core::TraceNode {
650 entity_type: node.entity.clone(),
651 entity_id: node.id().and_then(|v| match v {
652 Value::U64(n) => Some(*n),
653 Value::I64(n) => Some(*n as u64),
654 _ => None,
655 }),
656 comment: c.clone(),
657 },
658 });
659 let active_scope = current_scope.as_ref().or(parent_scope);
660
661 match node.operation {
662 GraphOperation::Upsert | GraphOperation::Create => {}
663 GraphOperation::Reference => {
664 return self
665 .validate_reference_node(
666 node,
667 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
668 )
669 .await;
670 }
671 GraphOperation::Remove => {
672 self.validate_remove_node(
673 &node,
674 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
675 )
676 .await?;
677 self.delete_graph_node(&node, parent_scope).await?;
678 return Ok(node);
679 }
680 }
681
682 let descriptor = self
683 .data_service
684 .metadata
685 .context
686 .require_entity(&node.entity)
687 .map_err(DataServiceError::Runtime)?;
688 let Some(id_property) = descriptor.id_property() else {
689 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
690 "entity {} has no id property for graph upsert",
691 node.entity
692 ))));
693 };
694 let Some(id) = node
695 .values
696 .get(&id_property.name)
697 .filter(|value| !is_unassigned_id_value(value))
698 .cloned()
699 else {
700 node.comment = None;
702 return self.insert_graph_node_scoped(node, active_scope).await;
703 };
704
705 if node.operation == GraphOperation::Create
706 || self
707 .fetch_graph_current_row_internal(
708 &node.entity,
709 &id_property.name,
710 &id,
711 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
712 )
713 .await?
714 .is_none()
715 {
716 node.comment = None;
717 return self.insert_graph_node_scoped(node, active_scope).await;
718 }
719
720 let mut one_relations = Vec::new();
721 let mut many_relations = Vec::new();
722 for (name, children) in std::mem::take(&mut node.relations) {
723 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
724 DataServiceError::Runtime(RuntimeError::MissingRelation {
725 entity: node.entity.clone(),
726 relation: name.clone(),
727 })
728 })?;
729 match relation.many {
730 true => many_relations.push((name, relation.clone(), children)),
731 false => one_relations.push((name, relation.clone(), children)),
732 }
733 }
734
735 for (name, relation, children) in one_relations {
736 if children.len() > 1 {
737 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
738 "relation {}.{} expects one child, got {}",
739 node.entity,
740 name,
741 children.len()
742 ))));
743 }
744 let mut saved_children = Vec::new();
745 for child in children {
746 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
747 let child_repo = self.scoped_data_service_internal(child.entity.clone());
748 let saved_child = child_repo
749 .upsert_graph_node_scoped(child, active_scope)
750 .await?;
751 if relation.attach {
752 let foreign_value = saved_child
753 .values
754 .get(&relation.foreign_key)
755 .cloned()
756 .ok_or_else(|| {
757 DataServiceError::Runtime(RuntimeError::Graph(format!(
758 "saved child {} missing foreign key {} for relation {}.{}",
759 relation.target_entity, relation.foreign_key, node.entity, name
760 )))
761 })?;
762 node.values
763 .insert(relation.local_key.clone(), foreign_value);
764 }
765 saved_children.push(saved_child);
766 }
767 node.relations.insert(name, saved_children);
768 }
769
770 let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
771 if !update.values.is_empty() {
772 let prepared_update = self
773 .prepare_update_command(&update)
774 .map_err(DataServiceError::Runtime)?;
775 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
776 self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
777 .await?;
778 for (field, value) in &prepared_update.values {
779 node.values.insert(field.clone(), value.clone());
780 }
781 if let Some(version_property) = descriptor.version_property() {
782 if let Some(expected_version) = prepared_update.expected_version {
783 node.values.insert(
784 version_property.name.clone(),
785 Value::I64(expected_version + 1),
786 );
787 }
788 }
789 }
790
791 for (name, relation, children) in many_relations {
792 let local_value =
793 node.values
794 .get(&relation.local_key)
795 .cloned()
796 .ok_or_else(|| {
797 DataServiceError::Runtime(RuntimeError::Graph(format!(
798 "parent {} missing local key {} for relation {}",
799 node.entity, relation.local_key, name
800 )))
801 })?;
802 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
803 let child_descriptor = self
804 .data_service
805 .metadata
806 .context
807 .require_entity(&relation.target_entity)
808 .map_err(DataServiceError::Runtime)?;
809 let child_id_property = child_descriptor.id_property().ok_or_else(|| {
810 DataServiceError::Runtime(RuntimeError::Graph(format!(
811 "entity {} has no id property",
812 relation.target_entity
813 )))
814 })?;
815
816 let mut seen = std::collections::BTreeSet::new();
817 let mut saved_children = Vec::new();
818 for mut child in children {
819 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
820 if relation.attach && child.operation != GraphOperation::Reference {
821 child
822 .values
823 .insert(relation.foreign_key.clone(), local_value.clone());
824 }
825 if let Some(child_id) = child
826 .values
827 .get(&child_id_property.name)
828 .filter(|value| !is_unassigned_id_value(value))
829 {
830 let key = graph_identity_key(child_id);
831 if !seen.insert(key.clone()) {
832 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
833 "duplicate child id {key} in relation {}.{}",
834 node.entity, name
835 ))));
836 }
837 }
838 saved_children.push(
839 child_repo
840 .upsert_graph_node_scoped(child, active_scope)
841 .await?,
842 );
843 }
844
845 node.relations.insert(name, saved_children);
846 }
847
848 Ok(node)
849 })
850 }
851
852 async fn validate_reference_node(
853 &self,
854 node: GraphNode,
855 trace_chain: Vec<teaql_core::TraceNode>,
856 ) -> Result<GraphNode, DataServiceError<E::Error>> {
857 if !node.relations.is_empty() {
858 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
859 "reference node {} cannot contain child relations",
860 node.entity
861 ))));
862 }
863 let descriptor = self
864 .data_service
865 .metadata
866 .context
867 .require_entity(&node.entity)
868 .map_err(DataServiceError::Runtime)?;
869 let id_property = descriptor.id_property().ok_or_else(|| {
870 DataServiceError::Runtime(RuntimeError::Graph(format!(
871 "entity {} has no id property for graph reference",
872 node.entity
873 )))
874 })?;
875 let id = node
876 .values
877 .get(&id_property.name)
878 .filter(|value| !is_unassigned_id_value(value))
879 .cloned()
880 .ok_or_else(|| {
881 DataServiceError::Runtime(RuntimeError::Graph(format!(
882 "reference node {} missing id property {}",
883 node.entity, id_property.name
884 )))
885 })?;
886
887 for field in node.values.keys() {
888 if field == &id_property.name {
889 continue;
890 }
891 if descriptor
892 .version_property()
893 .map(|property| field == &property.name)
894 .unwrap_or(false)
895 {
896 continue;
897 }
898 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
899 "reference node {} cannot carry mutable field {}",
900 node.entity, field
901 ))));
902 }
903
904 let current = self
905 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
906 .await?
907 .ok_or_else(|| {
908 DataServiceError::Runtime(RuntimeError::Graph(format!(
909 "reference node {}({}) does not exist",
910 node.entity,
911 graph_identity_key(&id)
912 )))
913 })?;
914
915 if let Some(version_property) = descriptor.version_property() {
916 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
917 if *existing_version < 0 {
918 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
919 "reference node {}({}) is deleted",
920 node.entity,
921 graph_identity_key(&id)
922 ))));
923 }
924 if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
925 {
926 if expected_version != existing_version {
927 println!(
928 "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
929 node.entity, expected_version, existing_version
930 );
931 return Err(DataServiceError::Runtime(
932 RuntimeError::OptimisticLockConflict {
933 entity: node.entity,
934 id: graph_identity_key(&id),
935 },
936 ));
937 }
938 }
939 }
940 }
941
942 Ok(GraphNode {
943 entity: node.entity,
944 values: current,
945 relations: BTreeMap::new(),
946 operation: GraphOperation::Reference,
947 comment: None,
948 dirty_fields: None,
949 original_values: None,
950 })
951 }
952
953 async fn validate_remove_node(
954 &self,
955 node: &GraphNode,
956 trace_chain: Vec<teaql_core::TraceNode>,
957 ) -> Result<(), DataServiceError<E::Error>> {
958 if !node.relations.is_empty() {
959 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
960 "remove node {} cannot contain child relations",
961 node.entity
962 ))));
963 }
964 let descriptor = self
965 .data_service
966 .metadata
967 .context
968 .require_entity(&node.entity)
969 .map_err(DataServiceError::Runtime)?;
970 let id_property = descriptor.id_property().ok_or_else(|| {
971 DataServiceError::Runtime(RuntimeError::Graph(format!(
972 "entity {} has no id property for graph remove",
973 node.entity
974 )))
975 })?;
976 let id = node
977 .values
978 .get(&id_property.name)
979 .filter(|value| !is_unassigned_id_value(value))
980 .cloned()
981 .ok_or_else(|| {
982 DataServiceError::Runtime(RuntimeError::Graph(format!(
983 "remove node {} missing id property {}",
984 node.entity, id_property.name
985 )))
986 })?;
987 let current = self
988 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
989 .await?
990 .ok_or_else(|| {
991 DataServiceError::Runtime(RuntimeError::Graph(format!(
992 "remove node {}({}) does not exist",
993 node.entity,
994 graph_identity_key(&id)
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 "remove node {}({}) is already deleted",
1002 node.entity,
1003 graph_identity_key(&id)
1004 ))));
1005 }
1006 }
1007 }
1008 Ok(())
1009 }
1010
1011 fn graph_node_from_record(
1012 &self,
1013 entity: &str,
1014 record: Record,
1015 ) -> Result<GraphNode, RuntimeError> {
1016 let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1017 let mut node = GraphNode::new(entity);
1018
1019 for (field, value) in record {
1020 if field == "_comment" {
1021 if let Value::Text(comment) = value {
1022 node.set_comment(comment);
1023 }
1024 continue;
1025 }
1026 if field == "_dirty_fields" {
1027 if let Value::List(fields) = value {
1028 let mut dirty = std::collections::BTreeSet::new();
1029 for f in fields {
1030 if let Value::Text(t) = f {
1031 dirty.insert(t);
1032 }
1033 }
1034 node.dirty_fields = Some(dirty);
1035 }
1036 continue;
1037 }
1038 if field == "_original_values" {
1039 if let Value::Object(orig) = value {
1040 node.original_values = Some(orig);
1041 }
1042 continue;
1043 }
1044 let Some(relation) = descriptor.relation_by_name(&field) else {
1045 node.values.insert(field, value);
1046 continue;
1047 };
1048
1049 match value {
1050 Value::Null => {
1051 node.relations.entry(field).or_default();
1052 }
1053 Value::Object(record) => {
1054 let child = self.graph_node_from_record(&relation.target_entity, record)?;
1055 node.relations.entry(field).or_default().push(child);
1056 }
1057 Value::List(values) => {
1058 let children = node.relations.entry(field.clone()).or_default();
1059 for value in values {
1060 let Value::Object(record) = value else {
1061 return Err(RuntimeError::Graph(format!(
1062 "relation {}.{} expects object children, got {:?}",
1063 entity, field, value
1064 )));
1065 };
1066 children
1067 .push(self.graph_node_from_record(&relation.target_entity, record)?);
1068 }
1069 }
1070 other => {
1071 return Err(RuntimeError::Graph(format!(
1072 "relation {}.{} expects object/list/null, got {:?}",
1073 entity, field, other
1074 )));
1075 }
1076 }
1077 }
1078
1079 Ok(node)
1080 }
1081
1082 fn graph_update_command(
1083 &self,
1084 node: &mut GraphNode,
1085 descriptor: &EntityDescriptor,
1086 id_property: &PropertyDescriptor,
1087 id: &Value,
1088 ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1089 crate::mark_record_status(&mut node.values, crate::CheckObjectStatus::Update);
1090 let check_result = self
1091 .data_service
1092 .metadata
1093 .context
1094 .check_and_fix_record(&node.entity, &mut node.values);
1095 crate::clear_record_status(&mut node.values);
1096 check_result.map_err(DataServiceError::Runtime)?;
1097
1098 let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1099 command.old_values = node.original_values.clone();
1100 if let Some(version_property) = descriptor.version_property() {
1101 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1102 command = command.expected_version(*version);
1103 }
1104 }
1105 for property in descriptor.properties.iter().filter(|property| {
1109 !property.is_id
1110 && !property.is_version
1111 && property.name != id_property.name
1112 && match &node.dirty_fields {
1113 Some(dirty) => dirty.contains(&property.name),
1114 None => node.values.contains_key(&property.name),
1115 }
1116 }) {
1117 if let Some(value) = node.values.get(&property.name) {
1118 command.values.insert(property.name.clone(), value.clone());
1119 }
1120 }
1121 Ok(command)
1122 }
1123
1124 fn delete_graph_node<'b, 's: 'b>(
1125 &'b self,
1126 node: &'b GraphNode,
1127 parent_scope: Option<&'s ScopedCommentNode<'s>>,
1128 ) -> std::pin::Pin<
1129 Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1130 > {
1131 Box::pin(async move {
1132 let descriptor = self
1133 .data_service
1134 .metadata
1135 .context
1136 .require_entity(&node.entity)
1137 .map_err(DataServiceError::Runtime)?;
1138 let id_property = descriptor.id_property().ok_or_else(|| {
1139 DataServiceError::Runtime(RuntimeError::Graph(format!(
1140 "entity {} has no id property for graph remove",
1141 node.entity
1142 )))
1143 })?;
1144 let id = node
1145 .values
1146 .get(&id_property.name)
1147 .filter(|value| !is_unassigned_id_value(value))
1148 .cloned()
1149 .ok_or_else(|| {
1150 DataServiceError::Runtime(RuntimeError::Graph(format!(
1151 "remove node {} missing id property {}",
1152 node.entity, id_property.name
1153 )))
1154 })?;
1155 let mut delete = DeleteCommand::new(node.entity.clone(), id);
1156 if let Some(version_property) = descriptor.version_property() {
1157 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1158 delete = delete.expected_version(*version);
1159 }
1160 }
1161
1162 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1164 parent: parent_scope,
1165 track: teaql_core::TraceNode {
1166 entity_type: node.entity.clone(),
1167 entity_id: node.id().and_then(|v| match v {
1168 Value::U64(n) => Some(*n),
1169 Value::I64(n) => Some(*n as u64),
1170 _ => None,
1171 }),
1172 comment: c.clone(),
1173 },
1174 });
1175 let active_scope = current_scope.as_ref().or(parent_scope);
1176 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1177
1178 self.delete_scoped_internal(&delete, lineage).await
1179 })
1180 }
1181
1182 async fn fetch_graph_children(
1183 &self,
1184 entity: &str,
1185 foreign_key: &str,
1186 parent_value: &Value,
1187 trace_chain: Vec<teaql_core::TraceNode>,
1188 ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
1189 let mut query =
1190 SelectQuery::new(entity).filter(Expr::eq(foreign_key, parent_value.clone()));
1191 query.trace_chain = trace_chain;
1192 self.scoped_data_service_internal(entity.to_owned())
1193 .fetch_all_internal(&query)
1194 .await
1195 }
1196 pub(crate) async fn fetch_graph_current_row_internal(
1197 &self,
1198 entity: &str,
1199 id_property: &str,
1200 id: &teaql_core::Value,
1201 trace_chain: Vec<teaql_core::TraceNode>,
1202 ) -> Result<Option<Record>, DataServiceError<E::Error>> {
1203 let mut query = teaql_core::SelectQuery::new(entity)
1204 .filter(teaql_core::Expr::eq(id_property, id.clone()));
1205 query.trace_chain = trace_chain;
1206 let mut rows = self
1207 .scoped_data_service_internal(entity.to_owned())
1208 .fetch_all_internal(&query)
1209 .await?;
1210 Ok(rows.pop())
1211 }
1212
1213 pub(crate) async fn execute_ledger_plan_internal(
1214 &self,
1215 root: crate::EntityRoot,
1216 ) -> Result<(), DataServiceError<E::Error>> {
1217 let comment = root.get_comment();
1218 let trace_chain = comment
1219 .map(|c| {
1220 vec![teaql_core::TraceNode {
1221 entity_type: self.entity.clone(),
1222 entity_id: None,
1223 comment: c,
1224 }]
1225 })
1226 .unwrap_or_default();
1227
1228 let deleted_keys = root.deleted_keys();
1229 let new_keys = root.new_keys();
1230 let change_set = root.current_change_set();
1231
1232 for key in deleted_keys.iter() {
1234 let id = key.id.clone();
1235 let mut cmd = teaql_core::DeleteCommand::new(&key.entity, id);
1236 if let Some(version) = root.get_original_version(key) {
1237 cmd = cmd.expected_version(version);
1238 }
1239 cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1240 self.delete_internal(&cmd).await?;
1241 }
1242
1243 let mut update_batches: std::collections::BTreeMap<
1245 (String, String),
1246 Vec<crate::EntityKey>,
1247 > = std::collections::BTreeMap::new();
1248 let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1249 std::collections::BTreeMap::new();
1250
1251 for (key, record) in change_set.changes() {
1252 if deleted_keys.contains(key) {
1253 continue;
1254 }
1255 let mut is_new = new_keys.contains(key);
1256
1257 if !is_new {
1258 let descriptor = self
1259 .data_service
1260 .metadata
1261 .context
1262 .require_entity(&key.entity)
1263 .map_err(DataServiceError::Runtime)?;
1264 let id_property = descriptor.id_property().ok_or_else(|| {
1265 DataServiceError::Runtime(RuntimeError::Graph(format!(
1266 "entity {} has no id property",
1267 key.entity
1268 )))
1269 })?;
1270 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1271 let current_row = self
1272 .fetch_graph_current_row_internal(
1273 &key.entity,
1274 &id_property.name,
1275 &key.id,
1276 my_trace,
1277 )
1278 .await?;
1279 if current_row.is_none() {
1280 is_new = true;
1281 }
1282 }
1283
1284 match is_new {
1285 true => {
1286 insert_batches
1287 .entry(key.entity.clone())
1288 .or_default()
1289 .push(key.clone());
1290 }
1291 false => {
1292 let mut fields: Vec<String> = record.keys().cloned().collect();
1293 fields.sort();
1294 let signature = fields.join(",");
1295 update_batches
1296 .entry((key.entity.clone(), signature))
1297 .or_default()
1298 .push(key.clone());
1299 }
1300 }
1301 }
1302
1303 let mut insert_order: Vec<String> = insert_batches.keys().cloned().collect();
1304 insert_order.sort();
1305 println!("execute_ledger_plan: insert_batches={:?}", insert_order);
1306
1307 for entity in insert_order {
1308 let keys = insert_batches.get(&entity).unwrap();
1309 let descriptor = self
1310 .data_service
1311 .metadata
1312 .context
1313 .require_entity(&entity)
1314 .map_err(DataServiceError::Runtime)?;
1315 let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.table_name);
1316 let mut traces = Vec::new();
1317 for key in keys {
1318 let record = change_set.changes().get(key).unwrap();
1319 let mut db_record = Record::new();
1320 db_record.insert("id".to_owned(), key.id.clone());
1321 for (field, value) in record {
1322 db_record.insert(field.clone(), value.clone());
1323 }
1324 crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1325 cmd.batch_values.push(db_record);
1326 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1327 traces.push(my_trace);
1328 }
1329 cmd.trace_chains = traces;
1330 self.execute_prepared_batch_insert(cmd).await?;
1331 }
1332
1333 let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1334 update_order.sort();
1335 println!("execute_ledger_plan: update_batches={:?}", update_order);
1336
1337 for signature in update_order {
1338 let keys = update_batches.get(&signature).unwrap();
1339 let descriptor = self
1340 .data_service
1341 .metadata
1342 .context
1343 .require_entity(&signature.0)
1344 .map_err(DataServiceError::Runtime)?;
1345 let mut update_fields: Vec<String> =
1346 signature.0.split(',').map(|s| s.to_string()).collect();
1347 let mut cmd =
1348 teaql_core::BatchUpdateCommand::new(&descriptor.table_name, update_fields);
1349 let mut traces = Vec::new();
1350 for key in keys {
1351 let record = change_set.changes().get(key).unwrap();
1352 let mut db_record = Record::new();
1353 db_record.insert("id".to_owned(), key.id.clone());
1354 for (field, value) in record {
1355 db_record.insert(field.clone(), value.clone());
1356 }
1357 crate::data_service::helpers::increment_version(
1358 &mut db_record,
1359 descriptor,
1360 root.get_original_version(key),
1361 );
1362 cmd.batch_values.push(db_record);
1363 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1364 traces.push(my_trace);
1365 }
1366 cmd.trace_chains = traces;
1367 self.execute_prepared_batch_update(cmd).await?;
1368 }
1369
1370 Ok(())
1371 }
1372}