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