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