1use std::collections::BTreeMap;
2use std::sync::Arc;
3
4use teaql_core::{
5 DeleteCommand, Entity, EntityDescriptor, Expr, InsertCommand, MutationValues,
6 PropertyDescriptor, 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.values.into(),
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.into());
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.into());
237 cmd.batch_ids.push(id);
238 cmd.batch_expected_versions.push(version);
239 cmd.batch_old_values.push(item.old_values.map(Into::into));
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_values(&descriptor.name, entity.into_values())?;
289 node.dirty_fields = dirty_fields;
290 node.original_values = original_values.map(Into::into);
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().into(),
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().into(),
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().into(),
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().into(),
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.into();
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 .map(Into::into)
618 .ok_or_else(|| {
619 DataServiceError::Runtime(RuntimeError::Graph(format!(
620 "persisted {} record could not be read back",
621 node.entity
622 )))
623 })?;
624 }
625 }
626
627 for (name, relation, children) in many_relations {
628 let local_value =
629 node.values
630 .get(&relation.local_key)
631 .cloned()
632 .ok_or_else(|| {
633 DataServiceError::Runtime(RuntimeError::Graph(format!(
634 "parent {} missing local key {} for relation {}",
635 node.entity, relation.local_key, name
636 )))
637 })?;
638 let mut saved_children = Vec::new();
639 for mut child in children {
640 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
641 if relation.attach {
642 child
643 .values
644 .insert(relation.foreign_key.clone(), local_value.clone());
645 }
646 let child_repo = self.scoped_data_service_internal(child.entity.clone());
647 saved_children.push(
648 child_repo
649 .insert_graph_node_scoped(child, active_scope)
650 .await?,
651 );
652 }
653 node.relations.insert(name, saved_children);
654 }
655
656 Ok(node)
657 })
658 }
659
660 fn upsert_graph_node_scoped<'b, 's: 'b>(
661 &'b self,
662 mut node: GraphNode,
663 parent_scope: Option<&'s ScopedCommentNode<'s>>,
664 ) -> std::pin::Pin<
665 Box<
666 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
667 + Send
668 + '_,
669 >,
670 > {
671 Box::pin(async move {
672 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
674 parent: parent_scope,
675 track: teaql_core::TraceNode {
676 entity_type: node.entity.clone(),
677 entity_id: node.id().and_then(|v| match v {
678 Value::U64(n) => Some(*n),
679 Value::I64(n) => Some(*n as u64),
680 _ => None,
681 }),
682 comment: c.clone(),
683 },
684 });
685 let active_scope = current_scope.as_ref().or(parent_scope);
686
687 match node.operation {
688 GraphOperation::Upsert | GraphOperation::Create => {}
689 GraphOperation::Reference => {
690 return self
691 .validate_reference_node(
692 node,
693 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
694 )
695 .await;
696 }
697 GraphOperation::Remove => {
698 self.validate_remove_node(
699 &node,
700 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
701 )
702 .await?;
703 self.delete_graph_node(&node, parent_scope).await?;
704 return Ok(node);
705 }
706 }
707
708 let descriptor = self
709 .data_service
710 .metadata
711 .context
712 .require_entity(&node.entity)
713 .map_err(DataServiceError::Runtime)?;
714 let Some(id_property) = descriptor.id_property() else {
715 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
716 "entity {} has no id property for graph upsert",
717 node.entity
718 ))));
719 };
720 let Some(id) = node
721 .values
722 .get(&id_property.name)
723 .filter(|value| !is_unassigned_id_value(value))
724 .cloned()
725 else {
726 node.comment = None;
728 return self.insert_graph_node_scoped(node, active_scope).await;
729 };
730
731 if node.operation == GraphOperation::Create
732 || self
733 .fetch_graph_current_row_internal(
734 &node.entity,
735 &id_property.name,
736 &id,
737 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
738 )
739 .await?
740 .is_none()
741 {
742 node.comment = None;
743 return self.insert_graph_node_scoped(node, active_scope).await;
744 }
745
746 let mut one_relations = Vec::new();
747 let mut many_relations = Vec::new();
748 for (name, children) in std::mem::take(&mut node.relations) {
749 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
750 DataServiceError::Runtime(RuntimeError::MissingRelation {
751 entity: node.entity.clone(),
752 relation: name.clone(),
753 })
754 })?;
755 match relation.many {
756 true => many_relations.push((name, relation.clone(), children)),
757 false => one_relations.push((name, relation.clone(), children)),
758 }
759 }
760
761 for (name, relation, children) in one_relations {
762 if children.len() > 1 {
763 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
764 "relation {}.{} expects one child, got {}",
765 node.entity,
766 name,
767 children.len()
768 ))));
769 }
770 let mut saved_children = Vec::new();
771 for child in children {
772 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
773 let child_repo = self.scoped_data_service_internal(child.entity.clone());
774 let saved_child = child_repo
775 .upsert_graph_node_scoped(child, active_scope)
776 .await?;
777 if relation.attach {
778 let foreign_value = saved_child
779 .values
780 .get(&relation.foreign_key)
781 .cloned()
782 .ok_or_else(|| {
783 DataServiceError::Runtime(RuntimeError::Graph(format!(
784 "saved child {} missing foreign key {} for relation {}.{}",
785 relation.target_entity, relation.foreign_key, node.entity, name
786 )))
787 })?;
788 node.values
789 .insert(relation.local_key.clone(), foreign_value);
790 }
791 saved_children.push(saved_child);
792 }
793 node.relations.insert(name, saved_children);
794 }
795
796 let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
797 if !update.values.is_empty() {
798 let prepared_update = self
799 .prepare_update_command(&update)
800 .map_err(DataServiceError::Runtime)?;
801 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
802 self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
803 .await?;
804 for (field, value) in &prepared_update.values {
805 node.values.insert(field.clone(), value.clone());
806 }
807 if let Some(version_property) = descriptor.version_property() {
808 if let Some(expected_version) = prepared_update.expected_version {
809 node.values.insert(
810 version_property.name.clone(),
811 Value::I64(expected_version + 1),
812 );
813 }
814 }
815 node.values = self
816 .fetch_graph_current_row_internal(
817 &node.entity,
818 &id_property.name,
819 &id,
820 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
821 )
822 .await?
823 .map(Into::into)
824 .ok_or_else(|| {
825 DataServiceError::Runtime(RuntimeError::Graph(format!(
826 "persisted {} record could not be read back",
827 node.entity
828 )))
829 })?;
830 }
831
832 for (name, relation, children) in many_relations {
833 let local_value =
834 node.values
835 .get(&relation.local_key)
836 .cloned()
837 .ok_or_else(|| {
838 DataServiceError::Runtime(RuntimeError::Graph(format!(
839 "parent {} missing local key {} for relation {}",
840 node.entity, relation.local_key, name
841 )))
842 })?;
843 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
844 let child_descriptor = self
845 .data_service
846 .metadata
847 .context
848 .require_entity(&relation.target_entity)
849 .map_err(DataServiceError::Runtime)?;
850 let child_id_property = child_descriptor.id_property().ok_or_else(|| {
851 DataServiceError::Runtime(RuntimeError::Graph(format!(
852 "entity {} has no id property",
853 relation.target_entity
854 )))
855 })?;
856
857 let mut seen = std::collections::BTreeSet::new();
858 let mut saved_children = Vec::new();
859 for mut child in children {
860 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
861 if relation.attach && child.operation != GraphOperation::Reference {
862 child
863 .values
864 .insert(relation.foreign_key.clone(), local_value.clone());
865 }
866 if let Some(child_id) = child
867 .values
868 .get(&child_id_property.name)
869 .filter(|value| !is_unassigned_id_value(value))
870 {
871 let key = graph_identity_key(child_id);
872 if !seen.insert(key.clone()) {
873 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
874 "duplicate child id {key} in relation {}.{}",
875 node.entity, name
876 ))));
877 }
878 }
879 saved_children.push(
880 child_repo
881 .upsert_graph_node_scoped(child, active_scope)
882 .await?,
883 );
884 }
885
886 node.relations.insert(name, saved_children);
887 }
888
889 Ok(node)
890 })
891 }
892
893 async fn validate_reference_node(
894 &self,
895 node: GraphNode,
896 trace_chain: Vec<teaql_core::TraceNode>,
897 ) -> Result<GraphNode, DataServiceError<E::Error>> {
898 if !node.relations.is_empty() {
899 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
900 "reference node {} cannot contain child relations",
901 node.entity
902 ))));
903 }
904 let descriptor = self
905 .data_service
906 .metadata
907 .context
908 .require_entity(&node.entity)
909 .map_err(DataServiceError::Runtime)?;
910 let id_property = descriptor.id_property().ok_or_else(|| {
911 DataServiceError::Runtime(RuntimeError::Graph(format!(
912 "entity {} has no id property for graph reference",
913 node.entity
914 )))
915 })?;
916 let id = node
917 .values
918 .get(&id_property.name)
919 .filter(|value| !is_unassigned_id_value(value))
920 .cloned()
921 .ok_or_else(|| {
922 DataServiceError::Runtime(RuntimeError::Graph(format!(
923 "reference node {} missing id property {}",
924 node.entity, id_property.name
925 )))
926 })?;
927
928 for field in node.values.keys() {
929 if field == &id_property.name {
930 continue;
931 }
932 if descriptor
933 .version_property()
934 .map(|property| field == &property.name)
935 .unwrap_or(false)
936 {
937 continue;
938 }
939 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
940 "reference node {} cannot carry mutable field {}",
941 node.entity, field
942 ))));
943 }
944
945 let current = self
946 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
947 .await?
948 .ok_or_else(|| {
949 DataServiceError::Runtime(RuntimeError::Graph(format!(
950 "reference node {}({}) does not exist",
951 node.entity,
952 graph_identity_key(&id)
953 )))
954 })?;
955
956 if let Some(version_property) = descriptor.version_property() {
957 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
958 if *existing_version < 0 {
959 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
960 "reference node {}({}) is deleted",
961 node.entity,
962 graph_identity_key(&id)
963 ))));
964 }
965 if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
966 {
967 if expected_version != existing_version {
968 println!(
969 "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
970 node.entity, expected_version, existing_version
971 );
972 return Err(DataServiceError::Runtime(
973 RuntimeError::OptimisticLockConflict {
974 entity: node.entity,
975 id: graph_identity_key(&id),
976 },
977 ));
978 }
979 }
980 }
981 }
982
983 Ok(GraphNode {
984 entity: node.entity,
985 values: current.into(),
986 relations: BTreeMap::new(),
987 operation: GraphOperation::Reference,
988 comment: None,
989 dirty_fields: None,
990 original_values: None,
991 })
992 }
993
994 async fn validate_remove_node(
995 &self,
996 node: &GraphNode,
997 trace_chain: Vec<teaql_core::TraceNode>,
998 ) -> Result<(), DataServiceError<E::Error>> {
999 if !node.relations.is_empty() {
1000 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1001 "remove node {} cannot contain child relations",
1002 node.entity
1003 ))));
1004 }
1005 let descriptor = self
1006 .data_service
1007 .metadata
1008 .context
1009 .require_entity(&node.entity)
1010 .map_err(DataServiceError::Runtime)?;
1011 let id_property = descriptor.id_property().ok_or_else(|| {
1012 DataServiceError::Runtime(RuntimeError::Graph(format!(
1013 "entity {} has no id property for graph remove",
1014 node.entity
1015 )))
1016 })?;
1017 let id = node
1018 .values
1019 .get(&id_property.name)
1020 .filter(|value| !is_unassigned_id_value(value))
1021 .cloned()
1022 .ok_or_else(|| {
1023 DataServiceError::Runtime(RuntimeError::Graph(format!(
1024 "remove node {} missing id property {}",
1025 node.entity, id_property.name
1026 )))
1027 })?;
1028 let current = self
1029 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
1030 .await?
1031 .ok_or_else(|| {
1032 DataServiceError::Runtime(RuntimeError::Graph(format!(
1033 "remove node {}({}) does not exist",
1034 node.entity,
1035 graph_identity_key(&id)
1036 )))
1037 })?;
1038 if let Some(version_property) = descriptor.version_property() {
1039 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
1040 if *existing_version < 0 {
1041 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1042 "remove node {}({}) is already deleted",
1043 node.entity,
1044 graph_identity_key(&id)
1045 ))));
1046 }
1047 }
1048 }
1049 Ok(())
1050 }
1051
1052 fn graph_node_from_values(
1053 &self,
1054 entity: &str,
1055 values: MutationValues,
1056 ) -> Result<GraphNode, RuntimeError> {
1057 let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1058 let mut node = GraphNode::new(entity);
1059
1060 for (field, value) in values {
1061 if field == "_comment" {
1062 if let Value::Text(comment) = value {
1063 node.set_comment(comment);
1064 }
1065 continue;
1066 }
1067 if field == "_dirty_fields" {
1068 if let Value::List(fields) = value {
1069 let mut dirty = std::collections::BTreeSet::new();
1070 for f in fields {
1071 if let Value::Text(t) = f {
1072 dirty.insert(t);
1073 }
1074 }
1075 node.dirty_fields = Some(dirty);
1076 }
1077 continue;
1078 }
1079 if field == "_original_values" {
1080 if let Value::Object(orig) = value {
1081 node.original_values = Some(orig.into());
1082 }
1083 continue;
1084 }
1085 if field == "_is_new" {
1086 if matches!(value, Value::Bool(true)) {
1087 node.operation = GraphOperation::Create;
1088 }
1089 continue;
1090 }
1091 if field == "_is_deleted" {
1092 if matches!(value, Value::Bool(true)) {
1093 node.operation = GraphOperation::Remove;
1094 }
1095 continue;
1096 }
1097 let Some(relation) = descriptor.relation_by_name(&field) else {
1098 node.values.insert(field, value);
1099 continue;
1100 };
1101
1102 match value {
1103 Value::Null => {
1104 node.relations.entry(field).or_default();
1105 }
1106 Value::Object(record) => {
1107 let child =
1108 self.graph_node_from_values(&relation.target_entity, record.into())?;
1109 node.relations.entry(field).or_default().push(child);
1110 }
1111 Value::List(values) => {
1112 let children = node.relations.entry(field.clone()).or_default();
1113 for value in values {
1114 let Value::Object(record) = value else {
1115 return Err(RuntimeError::Graph(format!(
1116 "relation {}.{} expects object children, got {:?}",
1117 entity, field, value
1118 )));
1119 };
1120 children.push(
1121 self.graph_node_from_values(&relation.target_entity, record.into())?,
1122 );
1123 }
1124 }
1125 other => {
1126 return Err(RuntimeError::Graph(format!(
1127 "relation {}.{} expects object/list/null, got {:?}",
1128 entity, field, other
1129 )));
1130 }
1131 }
1132 }
1133
1134 Ok(node)
1135 }
1136
1137 fn graph_update_command(
1138 &self,
1139 node: &mut GraphNode,
1140 descriptor: &EntityDescriptor,
1141 id_property: &PropertyDescriptor,
1142 id: &Value,
1143 ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1144 crate::mark_entity_status(&mut node.values, crate::CheckObjectStatus::Update);
1145 let check_result = self
1146 .data_service
1147 .metadata
1148 .context
1149 .check_and_fix_values(&node.entity, &mut node.values);
1150 crate::clear_entity_status(&mut node.values);
1151 check_result.map_err(DataServiceError::Runtime)?;
1152
1153 let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1154 command.old_values = node.original_values.clone().map(Into::into);
1155 if let Some(version_property) = descriptor.version_property() {
1156 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1157 command = command.expected_version(*version);
1158 }
1159 }
1160 for property in descriptor.properties.iter().filter(|property| {
1164 !property.is_id
1165 && !property.is_version
1166 && property.name != id_property.name
1167 && match &node.dirty_fields {
1168 Some(dirty) => dirty.contains(&property.name),
1169 None => node.values.contains_key(&property.name),
1170 }
1171 }) {
1172 if let Some(value) = node.values.get(&property.name) {
1173 command.values.insert(property.name.clone(), value.clone());
1174 }
1175 }
1176 Ok(command)
1177 }
1178
1179 fn delete_graph_node<'b, 's: 'b>(
1180 &'b self,
1181 node: &'b GraphNode,
1182 parent_scope: Option<&'s ScopedCommentNode<'s>>,
1183 ) -> std::pin::Pin<
1184 Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1185 > {
1186 Box::pin(async move {
1187 let descriptor = self
1188 .data_service
1189 .metadata
1190 .context
1191 .require_entity(&node.entity)
1192 .map_err(DataServiceError::Runtime)?;
1193 let id_property = descriptor.id_property().ok_or_else(|| {
1194 DataServiceError::Runtime(RuntimeError::Graph(format!(
1195 "entity {} has no id property for graph remove",
1196 node.entity
1197 )))
1198 })?;
1199 let id = node
1200 .values
1201 .get(&id_property.name)
1202 .filter(|value| !is_unassigned_id_value(value))
1203 .cloned()
1204 .ok_or_else(|| {
1205 DataServiceError::Runtime(RuntimeError::Graph(format!(
1206 "remove node {} missing id property {}",
1207 node.entity, id_property.name
1208 )))
1209 })?;
1210 let mut delete = DeleteCommand::new(node.entity.clone(), id);
1211 if let Some(version_property) = descriptor.version_property() {
1212 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1213 delete = delete.expected_version(*version);
1214 }
1215 }
1216
1217 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1219 parent: parent_scope,
1220 track: teaql_core::TraceNode {
1221 entity_type: node.entity.clone(),
1222 entity_id: node.id().and_then(|v| match v {
1223 Value::U64(n) => Some(*n),
1224 Value::I64(n) => Some(*n as u64),
1225 _ => None,
1226 }),
1227 comment: c.clone(),
1228 },
1229 });
1230 let active_scope = current_scope.as_ref().or(parent_scope);
1231 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1232
1233 self.delete_scoped_internal(&delete, lineage).await
1234 })
1235 }
1236
1237 pub(crate) async fn fetch_graph_current_row_internal(
1238 &self,
1239 entity: &str,
1240 id_property: &str,
1241 id: &teaql_core::Value,
1242 trace_chain: Vec<teaql_core::TraceNode>,
1243 ) -> Result<Option<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1244 let mut query = teaql_core::SelectQuery::new(entity)
1245 .filter(teaql_core::Expr::eq(id_property, id.clone()));
1246 query.trace_chain = trace_chain;
1247 let mut rows = self
1248 .scoped_data_service_internal(entity.to_owned())
1249 .fetch_all_internal(&query)
1250 .await?;
1251 Ok(rows.pop())
1252 }
1253
1254 pub(crate) async fn execute_ledger_plan_internal(
1255 &self,
1256 root: crate::EntityRuntimeState,
1257 ) -> Result<std::collections::BTreeMap<crate::EntityKey, Value>, DataServiceError<E::Error>>
1258 {
1259 let mut generated_ids = std::collections::BTreeMap::new();
1260 let comment = root.get_comment();
1261 let trace_chain = comment
1262 .map(|c| {
1263 vec![teaql_core::TraceNode {
1264 entity_type: self.entity.clone(),
1265 entity_id: None,
1266 comment: c,
1267 }]
1268 })
1269 .unwrap_or_default();
1270
1271 let deleted_keys = root.deleted_keys();
1272 let new_keys = root.new_keys();
1273 let change_set = root.current_change_set();
1274
1275 let mut checked_changes = std::collections::BTreeMap::new();
1282 for (key, record) in change_set.changes() {
1283 if deleted_keys.contains(key) {
1284 continue;
1285 }
1286 let mut checked: crate::EntityValues = record.clone().into();
1287 checked
1288 .entry("id".to_owned())
1289 .or_insert_with(|| key.id.clone());
1290 let status = if new_keys.contains(key) {
1291 crate::CheckObjectStatus::Create
1292 } else {
1293 crate::CheckObjectStatus::Update
1294 };
1295 crate::mark_entity_status(&mut checked, status);
1296 let result = self
1297 .data_service
1298 .metadata
1299 .context
1300 .check_and_fix_values(&key.entity, &mut checked);
1301 crate::clear_entity_status(&mut checked);
1302 result.map_err(DataServiceError::Runtime)?;
1303 for (field, value) in &checked {
1304 if record.get(field) != Some(value) {
1305 root.set(key.clone(), field, value.clone());
1306 }
1307 }
1308 checked_changes.insert(key.clone(), checked);
1309 }
1310
1311 for key in deleted_keys.iter() {
1313 let id = key.id.clone();
1314 let mut cmd = teaql_core::DeleteCommand::new(key.entity.as_ref(), id);
1315 if let Some(version) = root.get_original_version(key) {
1316 cmd = cmd.expected_version(version);
1317 }
1318 cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1319 self.delete_internal(&cmd).await?;
1320 }
1321
1322 let mut update_batches: std::collections::BTreeMap<
1324 (String, String),
1325 Vec<crate::EntityKey>,
1326 > = std::collections::BTreeMap::new();
1327 let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1328 std::collections::BTreeMap::new();
1329
1330 for (key, record) in &checked_changes {
1331 if deleted_keys.contains(key) {
1332 continue;
1333 }
1334 let mut is_new = new_keys.contains(key);
1335
1336 if !is_new {
1337 let descriptor = self
1338 .data_service
1339 .metadata
1340 .context
1341 .require_entity(&key.entity)
1342 .map_err(DataServiceError::Runtime)?;
1343 let id_property = descriptor.id_property().ok_or_else(|| {
1344 DataServiceError::Runtime(RuntimeError::Graph(format!(
1345 "entity {} has no id property",
1346 key.entity
1347 )))
1348 })?;
1349 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1350 let current_row = self
1351 .fetch_graph_current_row_internal(
1352 &key.entity,
1353 &id_property.name,
1354 &key.id,
1355 my_trace,
1356 )
1357 .await?;
1358 if current_row.is_none() {
1359 is_new = true;
1360 }
1361 }
1362
1363 match is_new {
1364 true => {
1365 insert_batches
1366 .entry(key.entity.to_string())
1367 .or_default()
1368 .push(key.clone());
1369 }
1370 false => {
1371 let mut fields: Vec<String> = record.keys().cloned().collect();
1372 fields.sort();
1373 let signature = fields.join(",");
1374 update_batches
1375 .entry((key.entity.to_string(), signature))
1376 .or_default()
1377 .push(key.clone());
1378 }
1379 }
1380 }
1381
1382 let mut insert_order: Vec<String> = insert_batches.keys().cloned().collect();
1383 insert_order.sort();
1384
1385 for entity in insert_order {
1386 let keys = insert_batches.get(&entity).unwrap();
1387 let descriptor = self
1388 .data_service
1389 .metadata
1390 .context
1391 .require_entity(&entity)
1392 .map_err(DataServiceError::Runtime)?;
1393 let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.name);
1394 let mut traces = Vec::new();
1395 for key in keys {
1396 let record = checked_changes.get(key).unwrap();
1397 let mut db_record = crate::EntityValues::new();
1398 let mut real_id = key.id.clone();
1399 if crate::data_service::helpers::is_unassigned_id_value(&real_id) {
1400 let gen_id = self
1401 .data_service
1402 .metadata
1403 .context
1404 .next_id(&entity)
1405 .map_err(DataServiceError::Runtime)?;
1406 real_id = Value::U64(gen_id);
1407 generated_ids.insert(key.clone(), real_id.clone());
1408 }
1409 db_record.insert("id".to_owned(), real_id);
1410 for (field, value) in record {
1411 if field == "id" {
1412 continue;
1413 }
1414 db_record.insert(field.clone(), value.clone());
1415 }
1416 crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1417 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, true);
1418 crate::mark_entity_status(&mut db_record, crate::CheckObjectStatus::Create);
1419 let check_result = self
1420 .data_service
1421 .metadata
1422 .context
1423 .check_and_fix_values(&entity, &mut db_record);
1424 crate::clear_entity_status(&mut db_record);
1425 check_result.map_err(DataServiceError::Runtime)?;
1426 cmd.batch_values.push(db_record.into());
1427 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1428 traces.push(my_trace);
1429 }
1430 cmd.trace_chains = traces;
1431 self.execute_prepared_batch_insert(cmd).await?;
1432 }
1433
1434 let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1435 update_order.sort();
1436
1437 for signature in update_order {
1438 let keys = update_batches.get(&signature).unwrap();
1439 let descriptor = self
1440 .data_service
1441 .metadata
1442 .context
1443 .require_entity(&signature.0)
1444 .map_err(DataServiceError::Runtime)?;
1445 let mut update_fields: Vec<String> =
1446 signature.1.split(',').map(|s| s.to_string()).collect();
1447 if descriptor
1448 .properties
1449 .iter()
1450 .any(|p| p.name == "update_time")
1451 && !update_fields.contains(&"update_time".to_owned())
1452 {
1453 update_fields.push("update_time".to_owned());
1454 }
1455 let mut cmd = teaql_core::BatchUpdateCommand::new(&descriptor.name, update_fields);
1456 let mut traces = Vec::new();
1457 for key in keys {
1458 let record = checked_changes.get(key).unwrap();
1459 let mut db_record = crate::EntityValues::new();
1460 db_record.insert("id".to_owned(), key.id.clone());
1461 for (field, value) in record {
1462 if field == "id" {
1463 continue;
1464 }
1465 db_record.insert(field.clone(), value.clone());
1466 }
1467 crate::data_service::helpers::increment_version(
1468 &mut db_record,
1469 descriptor,
1470 root.get_original_version(key),
1471 );
1472 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, false);
1473 crate::mark_entity_status(&mut db_record, crate::CheckObjectStatus::Update);
1474 let check_result = self
1475 .data_service
1476 .metadata
1477 .context
1478 .check_and_fix_values(&signature.0, &mut db_record);
1479 crate::clear_entity_status(&mut db_record);
1480 check_result.map_err(DataServiceError::Runtime)?;
1481 cmd.batch_values.push(db_record.into());
1482 cmd.batch_ids.push(key.id.clone());
1483 cmd.batch_expected_versions
1484 .push(root.get_original_version(key));
1485 cmd.batch_old_values.push(None); let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1487 traces.push(my_trace);
1488 }
1489 cmd.trace_chains = traces;
1490 self.execute_prepared_batch_update(cmd).await?;
1491 }
1492
1493 Ok(generated_ids)
1494 }
1495}