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