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