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 + teaql_data_service::MutationExecutor + Send + Sync,
37{
38 pub(crate) async fn save_graph_internal(
39 &self,
40 node: GraphNode,
41 ) -> Result<GraphNode, DataServiceError<E::Error>> {
42 if node.entity != self.entity {
43 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
44 "entity data service {} cannot save graph root {}",
45 self.entity, node.entity
46 ))));
47 }
48 let plan = self.plan_graph(node).await?;
49 self.execute_graph_plan_internal(plan).await
50 }
51
52 pub(crate) async fn save_entity_graph_from_internal(
53 &self,
54 graph: teaql_core::EntityGraph,
55 ) -> Result<GraphNode, DataServiceError<E::Error>> {
56 fn convert(node: teaql_core::EntityGraphNode) -> GraphNode {
57 let mut relations = BTreeMap::new();
58 for (rel_name, child) in node.children {
59 relations
60 .entry(rel_name)
61 .or_insert_with(Vec::new)
62 .push(convert(child));
63 }
64 GraphNode {
65 entity: node.entity_type,
66 values: node.values.into(),
67 relations,
68 operation: match node.operation {
69 teaql_core::EntityGraphOperation::Save => crate::GraphOperation::Upsert,
70 teaql_core::EntityGraphOperation::Delete => crate::GraphOperation::Remove,
71 },
72 comment: node.comment,
73 dirty_fields: None,
74 original_values: None,
75 }
76 }
77 self.save_graph_internal(convert(graph.root)).await
78 }
79
80 pub(crate) async fn save_entity_graph_internal<T>(
81 &self,
82 entity: T,
83 ) -> Result<GraphNode, DataServiceError<E::Error>>
84 where
85 T: Entity,
86 {
87 let node = self
88 .graph_node_from_entity(entity)
89 .map_err(DataServiceError::Runtime)?;
90 self.save_graph_internal(node).await
91 }
92
93 pub(crate) async fn save_entity_internal<T>(
94 &self,
95 entity: T,
96 status: EntityStatus,
97 ) -> Result<GraphNode, DataServiceError<E::Error>>
98 where
99 T: Entity,
100 {
101 if !status.need_persist() {
102 return Ok(GraphNode::new(&self.entity));
103 }
104 if status.is_deleted() {
105 let mut node = self
106 .graph_node_from_entity(entity)
107 .map_err(DataServiceError::Runtime)?;
108 node.operation = GraphOperation::Remove;
109 node.relations.clear();
110 return self.save_graph_internal(node).await;
111 }
112 self.save_entity_graph_internal(entity).await
113 }
114 pub(crate) async fn save_entity_with_comment_internal<T>(
115 &self,
116 entity: T,
117 status: EntityStatus,
118 comment: impl Into<String>,
119 ) -> Result<GraphNode, DataServiceError<E::Error>>
120 where
121 T: Entity,
122 {
123 if status.is_deleted() {
124 let mut node = self
125 .graph_node_from_entity(entity)
126 .map_err(DataServiceError::Runtime)?;
127 node.operation = GraphOperation::Remove;
128 node.relations.clear();
129 node.set_comment(comment);
130 return self.save_graph_internal(node).await;
131 }
132 self.save_entity_graph_with_comment_internal(entity, comment)
133 .await
134 }
135 pub(crate) async fn save_entity_graph_with_comment_internal<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_internal(node).await
148 }
149
150 pub(crate) async fn create_entity_graph_with_comment_internal<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_internal(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(crate) async fn execute_graph_plan_internal(
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.into());
210 cmd.trace_chains
211 .push(recover_trace_or_default(&item.scope_token));
212 }
213 self.execute_prepared_batch_insert(cmd).await?;
214 }
215 GraphMutationKind::Update => {
216 if batch.update_fields.is_empty() {
217 continue;
218 }
219 let mut cmd =
220 teaql_core::BatchUpdateCommand::new(&batch.entity, batch.update_fields);
221 for item in batch.items {
222 let id = item.values.get("id").cloned().ok_or_else(|| {
223 DataServiceError::Runtime(RuntimeError::Graph(format!(
224 "update item in batch missing id for {}",
225 batch.entity
226 )))
227 })?;
228 let version = item.values.get("version").and_then(|v| match v {
229 teaql_core::Value::I64(n) => Some(*n),
230 _ => None,
231 });
232 cmd.batch_values.push(item.values.into());
233 cmd.batch_ids.push(id);
234 cmd.batch_expected_versions.push(version);
235 cmd.batch_old_values.push(item.old_values.map(Into::into));
236 cmd.trace_chains
237 .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_internal(&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_values(&descriptor.name, entity.into_values())?;
285 node.dirty_fields = dirty_fields;
286 node.original_values = original_values.map(Into::into);
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().into(),
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().into(),
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_internal(
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 crate::data_service::helpers::ensure_timestamps(&mut node.values, descriptor, true);
410 } else {
411 crate::data_service::helpers::ensure_timestamps(
412 &mut node.values,
413 descriptor,
414 false,
415 );
416 }
417 let update_fields = is_update
418 .then(|| {
419 let mut excluded = Vec::new();
420 if let Some(id_property) = id_property.as_ref() {
421 excluded.push(id_property.name.clone());
422 }
423 if let Some(version_property) = descriptor.version_property() {
424 excluded.push(version_property.name.clone());
425 }
426 let mut fields = sorted_update_fields(&node.values, excluded);
427 if let Some(dirty) = &node.dirty_fields {
428 fields.retain(|f| dirty.contains(f));
429 }
430 fields
431 })
432 .unwrap_or_default();
433
434 let current_token = node
437 .comment
438 .as_ref()
439 .map(|c| {
440 Arc::new(TraceScopeToken {
441 parent: parent_token.clone(),
442 track: teaql_core::TraceNode {
443 entity_type: node.entity.clone(),
444 entity_id: node.id().and_then(|v| match v {
445 Value::U64(n) => Some(*n),
446 Value::I64(n) => Some(*n as u64),
447 _ => None,
448 }),
449 comment: c.clone(),
450 },
451 node_index: plan.next_item_index,
452 })
453 })
454 .or_else(|| parent_token.clone());
455
456 plan.push(
457 node.entity.clone(),
458 GraphMutationKind::for_update(is_update),
459 node.values.clone().into(),
460 update_fields,
461 current_token.clone(),
462 node.original_values.clone(),
463 );
464
465 for (name, children) in &mut node.relations {
466 let relation = descriptor.relation_by_name(name).ok_or_else(|| {
467 DataServiceError::Runtime(RuntimeError::MissingRelation {
468 entity: node.entity.clone(),
469 relation: name.clone(),
470 })
471 })?;
472 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
473 for child in children {
474 ensure_relation_target(&node.entity, name, &relation.target_entity, child)?;
475 child_repo
476 .collect_graph_plan(
477 child,
478 plan,
479 active_scope,
480 current_token.clone(),
481 is_create_op,
482 )
483 .await?;
484 }
485 }
486 Ok(())
487 })
488 }
489
490 fn insert_graph_node_scoped<'b, 's: 'b>(
491 &'b self,
492 mut node: GraphNode,
493 parent_scope: Option<&'s ScopedCommentNode<'s>>,
494 ) -> std::pin::Pin<
495 Box<
496 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
497 + Send
498 + '_,
499 >,
500 > {
501 Box::pin(async move {
502 match node.operation {
503 GraphOperation::Upsert | GraphOperation::Create => {}
504 GraphOperation::Reference => {
505 return self
506 .validate_reference_node(
507 node,
508 parent_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
509 )
510 .await;
511 }
512 GraphOperation::Remove => {
513 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
514 "create graph cannot remove node {}",
515 node.entity
516 ))));
517 }
518 }
519
520 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
522 parent: parent_scope,
523 track: teaql_core::TraceNode {
524 entity_type: node.entity.clone(),
525 entity_id: node.id().and_then(|v| match v {
526 Value::U64(n) => Some(*n),
527 Value::I64(n) => Some(*n as u64),
528 _ => None,
529 }),
530 comment: c.clone(),
531 },
532 });
533 let active_scope = current_scope.as_ref().or(parent_scope);
534
535 let descriptor = self
536 .data_service
537 .metadata
538 .context
539 .require_entity(&node.entity)
540 .map_err(DataServiceError::Runtime)?;
541
542 let mut one_relations = Vec::new();
543 let mut many_relations = Vec::new();
544 for (name, children) in std::mem::take(&mut node.relations) {
545 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
546 DataServiceError::Runtime(RuntimeError::MissingRelation {
547 entity: node.entity.clone(),
548 relation: name.clone(),
549 })
550 })?;
551 match relation.many {
552 true => many_relations.push((name, relation.clone(), children)),
553 false => one_relations.push((name, relation.clone(), children)),
554 }
555 }
556
557 for (name, relation, children) in one_relations {
558 if children.len() > 1 {
559 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
560 "relation {}.{} expects one child, got {}",
561 node.entity,
562 name,
563 children.len()
564 ))));
565 }
566 let mut saved_children = Vec::new();
567 for child in children {
568 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
569 let child_repo = self.scoped_data_service_internal(child.entity.clone());
570 let saved_child = child_repo
571 .insert_graph_node_scoped(child, active_scope)
572 .await?;
573 if relation.attach {
574 let foreign_value = saved_child
575 .values
576 .get(&relation.foreign_key)
577 .cloned()
578 .ok_or_else(|| {
579 DataServiceError::Runtime(RuntimeError::Graph(format!(
580 "saved child {} missing foreign key {} for relation {}.{}",
581 relation.target_entity, relation.foreign_key, node.entity, name
582 )))
583 })?;
584 node.values
585 .insert(relation.local_key.clone(), foreign_value);
586 }
587 saved_children.push(saved_child);
588 }
589 node.relations.insert(name, saved_children);
590 }
591
592 let command = self
593 .prepare_insert_command(&InsertCommand {
594 entity: node.entity.clone(),
595 values: node.values.clone().into(),
596 trace_chain: Vec::new(),
597 })
598 .map_err(DataServiceError::Runtime)?;
599 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
600 self.execute_prepared_insert_with_comment(command.clone(), lineage)
601 .await?;
602 node.values = command.values.into();
603 if let Some(id_property) = descriptor.id_property() {
604 if let Some(id) = node.values.get(&id_property.name).cloned() {
605 node.values = self
606 .fetch_graph_current_row_internal(
607 &node.entity,
608 &id_property.name,
609 &id,
610 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
611 )
612 .await?
613 .map(Into::into)
614 .ok_or_else(|| {
615 DataServiceError::Runtime(RuntimeError::Graph(format!(
616 "persisted {} record could not be read back",
617 node.entity
618 )))
619 })?;
620 }
621 }
622
623 for (name, relation, children) in many_relations {
624 let local_value =
625 node.values
626 .get(&relation.local_key)
627 .cloned()
628 .ok_or_else(|| {
629 DataServiceError::Runtime(RuntimeError::Graph(format!(
630 "parent {} missing local key {} for relation {}",
631 node.entity, relation.local_key, name
632 )))
633 })?;
634 let mut saved_children = Vec::new();
635 for mut child in children {
636 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
637 if relation.attach {
638 child
639 .values
640 .insert(relation.foreign_key.clone(), local_value.clone());
641 }
642 let child_repo = self.scoped_data_service_internal(child.entity.clone());
643 saved_children.push(
644 child_repo
645 .insert_graph_node_scoped(child, active_scope)
646 .await?,
647 );
648 }
649 node.relations.insert(name, saved_children);
650 }
651
652 Ok(node)
653 })
654 }
655
656 fn upsert_graph_node_scoped<'b, 's: 'b>(
657 &'b self,
658 mut node: GraphNode,
659 parent_scope: Option<&'s ScopedCommentNode<'s>>,
660 ) -> std::pin::Pin<
661 Box<
662 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
663 + Send
664 + '_,
665 >,
666 > {
667 Box::pin(async move {
668 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
670 parent: parent_scope,
671 track: teaql_core::TraceNode {
672 entity_type: node.entity.clone(),
673 entity_id: node.id().and_then(|v| match v {
674 Value::U64(n) => Some(*n),
675 Value::I64(n) => Some(*n as u64),
676 _ => None,
677 }),
678 comment: c.clone(),
679 },
680 });
681 let active_scope = current_scope.as_ref().or(parent_scope);
682
683 match node.operation {
684 GraphOperation::Upsert | GraphOperation::Create => {}
685 GraphOperation::Reference => {
686 return self
687 .validate_reference_node(
688 node,
689 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
690 )
691 .await;
692 }
693 GraphOperation::Remove => {
694 self.validate_remove_node(
695 &node,
696 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
697 )
698 .await?;
699 self.delete_graph_node(&node, parent_scope).await?;
700 return Ok(node);
701 }
702 }
703
704 let descriptor = self
705 .data_service
706 .metadata
707 .context
708 .require_entity(&node.entity)
709 .map_err(DataServiceError::Runtime)?;
710 let Some(id_property) = descriptor.id_property() else {
711 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
712 "entity {} has no id property for graph upsert",
713 node.entity
714 ))));
715 };
716 let Some(id) = node
717 .values
718 .get(&id_property.name)
719 .filter(|value| !is_unassigned_id_value(value))
720 .cloned()
721 else {
722 node.comment = None;
724 return self.insert_graph_node_scoped(node, active_scope).await;
725 };
726
727 if node.operation == GraphOperation::Create
728 || self
729 .fetch_graph_current_row_internal(
730 &node.entity,
731 &id_property.name,
732 &id,
733 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
734 )
735 .await?
736 .is_none()
737 {
738 node.comment = None;
739 return self.insert_graph_node_scoped(node, active_scope).await;
740 }
741
742 let mut one_relations = Vec::new();
743 let mut many_relations = Vec::new();
744 for (name, children) in std::mem::take(&mut node.relations) {
745 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
746 DataServiceError::Runtime(RuntimeError::MissingRelation {
747 entity: node.entity.clone(),
748 relation: name.clone(),
749 })
750 })?;
751 match relation.many {
752 true => many_relations.push((name, relation.clone(), children)),
753 false => one_relations.push((name, relation.clone(), children)),
754 }
755 }
756
757 for (name, relation, children) in one_relations {
758 if children.len() > 1 {
759 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
760 "relation {}.{} expects one child, got {}",
761 node.entity,
762 name,
763 children.len()
764 ))));
765 }
766 let mut saved_children = Vec::new();
767 for child in children {
768 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
769 let child_repo = self.scoped_data_service_internal(child.entity.clone());
770 let saved_child = child_repo
771 .upsert_graph_node_scoped(child, active_scope)
772 .await?;
773 if relation.attach {
774 let foreign_value = saved_child
775 .values
776 .get(&relation.foreign_key)
777 .cloned()
778 .ok_or_else(|| {
779 DataServiceError::Runtime(RuntimeError::Graph(format!(
780 "saved child {} missing foreign key {} for relation {}.{}",
781 relation.target_entity, relation.foreign_key, node.entity, name
782 )))
783 })?;
784 node.values
785 .insert(relation.local_key.clone(), foreign_value);
786 }
787 saved_children.push(saved_child);
788 }
789 node.relations.insert(name, saved_children);
790 }
791
792 let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
793 if !update.values.is_empty() {
794 let prepared_update = self
795 .prepare_update_command(&update)
796 .map_err(DataServiceError::Runtime)?;
797 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
798 self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
799 .await?;
800 for (field, value) in &prepared_update.values {
801 node.values.insert(field.clone(), value.clone());
802 }
803 if let Some(version_property) = descriptor.version_property() {
804 if let Some(expected_version) = prepared_update.expected_version {
805 node.values.insert(
806 version_property.name.clone(),
807 Value::I64(expected_version + 1),
808 );
809 }
810 }
811 node.values = self
812 .fetch_graph_current_row_internal(
813 &node.entity,
814 &id_property.name,
815 &id,
816 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
817 )
818 .await?
819 .map(Into::into)
820 .ok_or_else(|| {
821 DataServiceError::Runtime(RuntimeError::Graph(format!(
822 "persisted {} record could not be read back",
823 node.entity
824 )))
825 })?;
826 }
827
828 for (name, relation, children) in many_relations {
829 let local_value =
830 node.values
831 .get(&relation.local_key)
832 .cloned()
833 .ok_or_else(|| {
834 DataServiceError::Runtime(RuntimeError::Graph(format!(
835 "parent {} missing local key {} for relation {}",
836 node.entity, relation.local_key, name
837 )))
838 })?;
839 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
840 let child_descriptor = self
841 .data_service
842 .metadata
843 .context
844 .require_entity(&relation.target_entity)
845 .map_err(DataServiceError::Runtime)?;
846 let child_id_property = child_descriptor.id_property().ok_or_else(|| {
847 DataServiceError::Runtime(RuntimeError::Graph(format!(
848 "entity {} has no id property",
849 relation.target_entity
850 )))
851 })?;
852
853 let mut seen = std::collections::BTreeSet::new();
854 let mut saved_children = Vec::new();
855 for mut child in children {
856 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
857 if relation.attach && child.operation != GraphOperation::Reference {
858 child
859 .values
860 .insert(relation.foreign_key.clone(), local_value.clone());
861 }
862 if let Some(child_id) = child
863 .values
864 .get(&child_id_property.name)
865 .filter(|value| !is_unassigned_id_value(value))
866 {
867 let key = graph_identity_key(child_id);
868 if !seen.insert(key.clone()) {
869 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
870 "duplicate child id {key} in relation {}.{}",
871 node.entity, name
872 ))));
873 }
874 }
875 saved_children.push(
876 child_repo
877 .upsert_graph_node_scoped(child, active_scope)
878 .await?,
879 );
880 }
881
882 node.relations.insert(name, saved_children);
883 }
884
885 Ok(node)
886 })
887 }
888
889 async fn validate_reference_node(
890 &self,
891 node: GraphNode,
892 trace_chain: Vec<teaql_core::TraceNode>,
893 ) -> Result<GraphNode, DataServiceError<E::Error>> {
894 if !node.relations.is_empty() {
895 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
896 "reference node {} cannot contain child relations",
897 node.entity
898 ))));
899 }
900 let descriptor = self
901 .data_service
902 .metadata
903 .context
904 .require_entity(&node.entity)
905 .map_err(DataServiceError::Runtime)?;
906 let id_property = descriptor.id_property().ok_or_else(|| {
907 DataServiceError::Runtime(RuntimeError::Graph(format!(
908 "entity {} has no id property for graph reference",
909 node.entity
910 )))
911 })?;
912 let id = node
913 .values
914 .get(&id_property.name)
915 .filter(|value| !is_unassigned_id_value(value))
916 .cloned()
917 .ok_or_else(|| {
918 DataServiceError::Runtime(RuntimeError::Graph(format!(
919 "reference node {} missing id property {}",
920 node.entity, id_property.name
921 )))
922 })?;
923
924 for field in node.values.keys() {
925 if field == &id_property.name {
926 continue;
927 }
928 if descriptor
929 .version_property()
930 .map(|property| field == &property.name)
931 .unwrap_or(false)
932 {
933 continue;
934 }
935 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
936 "reference node {} cannot carry mutable field {}",
937 node.entity, field
938 ))));
939 }
940
941 let current = self
942 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
943 .await?
944 .ok_or_else(|| {
945 DataServiceError::Runtime(RuntimeError::Graph(format!(
946 "reference node {}({}) does not exist",
947 node.entity,
948 graph_identity_key(&id)
949 )))
950 })?;
951
952 if let Some(version_property) = descriptor.version_property() {
953 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
954 if *existing_version < 0 {
955 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
956 "reference node {}({}) is deleted",
957 node.entity,
958 graph_identity_key(&id)
959 ))));
960 }
961 if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
962 {
963 if expected_version != existing_version {
964 println!(
965 "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
966 node.entity, expected_version, existing_version
967 );
968 return Err(DataServiceError::Runtime(
969 RuntimeError::OptimisticLockConflict {
970 entity: node.entity,
971 id: graph_identity_key(&id),
972 },
973 ));
974 }
975 }
976 }
977 }
978
979 Ok(GraphNode {
980 entity: node.entity,
981 values: current.into(),
982 relations: BTreeMap::new(),
983 operation: GraphOperation::Reference,
984 comment: None,
985 dirty_fields: None,
986 original_values: None,
987 })
988 }
989
990 async fn validate_remove_node(
991 &self,
992 node: &GraphNode,
993 trace_chain: Vec<teaql_core::TraceNode>,
994 ) -> Result<(), DataServiceError<E::Error>> {
995 if !node.relations.is_empty() {
996 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
997 "remove node {} cannot contain child relations",
998 node.entity
999 ))));
1000 }
1001 let descriptor = self
1002 .data_service
1003 .metadata
1004 .context
1005 .require_entity(&node.entity)
1006 .map_err(DataServiceError::Runtime)?;
1007 let id_property = descriptor.id_property().ok_or_else(|| {
1008 DataServiceError::Runtime(RuntimeError::Graph(format!(
1009 "entity {} has no id property for graph remove",
1010 node.entity
1011 )))
1012 })?;
1013 let id = node
1014 .values
1015 .get(&id_property.name)
1016 .filter(|value| !is_unassigned_id_value(value))
1017 .cloned()
1018 .ok_or_else(|| {
1019 DataServiceError::Runtime(RuntimeError::Graph(format!(
1020 "remove node {} missing id property {}",
1021 node.entity, id_property.name
1022 )))
1023 })?;
1024 let current = self
1025 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
1026 .await?
1027 .ok_or_else(|| {
1028 DataServiceError::Runtime(RuntimeError::Graph(format!(
1029 "remove node {}({}) does not exist",
1030 node.entity,
1031 graph_identity_key(&id)
1032 )))
1033 })?;
1034 if let Some(version_property) = descriptor.version_property() {
1035 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
1036 if *existing_version < 0 {
1037 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1038 "remove node {}({}) is already deleted",
1039 node.entity,
1040 graph_identity_key(&id)
1041 ))));
1042 }
1043 }
1044 }
1045 Ok(())
1046 }
1047
1048 fn graph_node_from_values(
1049 &self,
1050 entity: &str,
1051 values: MutationValues,
1052 ) -> Result<GraphNode, RuntimeError> {
1053 let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1054 let mut node = GraphNode::new(entity);
1055
1056 for (field, value) in values {
1057 if field == "_comment" {
1058 if let Value::Text(comment) = value {
1059 node.set_comment(comment);
1060 }
1061 continue;
1062 }
1063 if field == "_dirty_fields" {
1064 if let Value::List(fields) = value {
1065 let mut dirty = std::collections::BTreeSet::new();
1066 for f in fields {
1067 if let Value::Text(t) = f {
1068 dirty.insert(t);
1069 }
1070 }
1071 node.dirty_fields = Some(dirty);
1072 }
1073 continue;
1074 }
1075 if field == "_original_values" {
1076 if let Value::Object(orig) = value {
1077 node.original_values = Some(orig.into());
1078 }
1079 continue;
1080 }
1081 if field == "_is_new" {
1082 if matches!(value, Value::Bool(true)) {
1083 node.operation = GraphOperation::Create;
1084 }
1085 continue;
1086 }
1087 if field == "_is_deleted" {
1088 if matches!(value, Value::Bool(true)) {
1089 node.operation = GraphOperation::Remove;
1090 }
1091 continue;
1092 }
1093 let Some(relation) = descriptor.relation_by_name(&field) else {
1094 node.values.insert(field, value);
1095 continue;
1096 };
1097
1098 match value {
1099 Value::Null => {
1100 node.relations.entry(field).or_default();
1101 }
1102 Value::Object(record) => {
1103 let child =
1104 self.graph_node_from_values(&relation.target_entity, record.into())?;
1105 node.relations.entry(field).or_default().push(child);
1106 }
1107 Value::List(values) => {
1108 let children = node.relations.entry(field.clone()).or_default();
1109 for value in values {
1110 let Value::Object(record) = value else {
1111 return Err(RuntimeError::Graph(format!(
1112 "relation {}.{} expects object children, got {:?}",
1113 entity, field, value
1114 )));
1115 };
1116 children.push(
1117 self.graph_node_from_values(&relation.target_entity, record.into())?,
1118 );
1119 }
1120 }
1121 other => {
1122 return Err(RuntimeError::Graph(format!(
1123 "relation {}.{} expects object/list/null, got {:?}",
1124 entity, field, other
1125 )));
1126 }
1127 }
1128 }
1129
1130 Ok(node)
1131 }
1132
1133 fn graph_update_command(
1134 &self,
1135 node: &mut GraphNode,
1136 descriptor: &EntityDescriptor,
1137 id_property: &PropertyDescriptor,
1138 id: &Value,
1139 ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1140 crate::mark_entity_status(&mut node.values, crate::CheckObjectStatus::Update);
1141 let check_result = self
1142 .data_service
1143 .metadata
1144 .context
1145 .check_and_fix_values(&node.entity, &mut node.values);
1146 crate::clear_entity_status(&mut node.values);
1147 check_result.map_err(DataServiceError::Runtime)?;
1148
1149 let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1150 command.old_values = node.original_values.clone().map(Into::into);
1151 if let Some(version_property) = descriptor.version_property() {
1152 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1153 command = command.expected_version(*version);
1154 }
1155 }
1156 for property in descriptor.properties.iter().filter(|property| {
1160 !property.is_id
1161 && !property.is_version
1162 && property.name != id_property.name
1163 && match &node.dirty_fields {
1164 Some(dirty) => dirty.contains(&property.name),
1165 None => node.values.contains_key(&property.name),
1166 }
1167 }) {
1168 if let Some(value) = node.values.get(&property.name) {
1169 command.values.insert(property.name.clone(), value.clone());
1170 }
1171 }
1172 Ok(command)
1173 }
1174
1175 fn delete_graph_node<'b, 's: 'b>(
1176 &'b self,
1177 node: &'b GraphNode,
1178 parent_scope: Option<&'s ScopedCommentNode<'s>>,
1179 ) -> std::pin::Pin<
1180 Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1181 > {
1182 Box::pin(async move {
1183 let descriptor = self
1184 .data_service
1185 .metadata
1186 .context
1187 .require_entity(&node.entity)
1188 .map_err(DataServiceError::Runtime)?;
1189 let id_property = descriptor.id_property().ok_or_else(|| {
1190 DataServiceError::Runtime(RuntimeError::Graph(format!(
1191 "entity {} has no id property for graph remove",
1192 node.entity
1193 )))
1194 })?;
1195 let id = node
1196 .values
1197 .get(&id_property.name)
1198 .filter(|value| !is_unassigned_id_value(value))
1199 .cloned()
1200 .ok_or_else(|| {
1201 DataServiceError::Runtime(RuntimeError::Graph(format!(
1202 "remove node {} missing id property {}",
1203 node.entity, id_property.name
1204 )))
1205 })?;
1206 let mut delete = DeleteCommand::new(node.entity.clone(), id);
1207 if let Some(version_property) = descriptor.version_property() {
1208 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1209 delete = delete.expected_version(*version);
1210 }
1211 }
1212
1213 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1215 parent: parent_scope,
1216 track: teaql_core::TraceNode {
1217 entity_type: node.entity.clone(),
1218 entity_id: node.id().and_then(|v| match v {
1219 Value::U64(n) => Some(*n),
1220 Value::I64(n) => Some(*n as u64),
1221 _ => None,
1222 }),
1223 comment: c.clone(),
1224 },
1225 });
1226 let active_scope = current_scope.as_ref().or(parent_scope);
1227 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1228
1229 self.delete_scoped_internal(&delete, lineage).await
1230 })
1231 }
1232
1233 pub(crate) async fn fetch_graph_current_row_internal(
1234 &self,
1235 entity: &str,
1236 id_property: &str,
1237 id: &teaql_core::Value,
1238 trace_chain: Vec<teaql_core::TraceNode>,
1239 ) -> Result<Option<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1240 let mut query = teaql_core::SelectQuery::new(entity)
1241 .filter(teaql_core::Expr::eq(id_property, id.clone()));
1242 query.trace_chain = trace_chain;
1243 let mut rows = self
1244 .scoped_data_service_internal(entity.to_owned())
1245 .fetch_all_internal(&query)
1246 .await?;
1247 Ok(rows.pop())
1248 }
1249
1250 pub(crate) async fn execute_ledger_plan_internal(
1251 &self,
1252 root: crate::EntityRuntimeState,
1253 ) -> Result<std::collections::BTreeMap<crate::EntityKey, Value>, DataServiceError<E::Error>>
1254 {
1255 let mut generated_ids = std::collections::BTreeMap::new();
1256 let comment = root.get_comment();
1257 let trace_chain = comment
1258 .map(|c| {
1259 vec![teaql_core::TraceNode {
1260 entity_type: self.entity.clone(),
1261 entity_id: None,
1262 comment: c,
1263 }]
1264 })
1265 .unwrap_or_default();
1266
1267 let deleted_keys = root.deleted_keys();
1268 let new_keys = root.new_keys();
1269 let change_set = root.current_change_set();
1270
1271 let mut checked_changes = std::collections::BTreeMap::new();
1278 for (key, record) in change_set.changes() {
1279 if deleted_keys.contains(key) {
1280 continue;
1281 }
1282 let mut checked: crate::EntityValues = record.clone().into();
1283 checked
1284 .entry("id".to_owned())
1285 .or_insert_with(|| key.id.clone());
1286 checked_changes.insert(key.clone(), checked);
1287 }
1288
1289 for key in deleted_keys.iter() {
1291 let id = key.id.clone();
1292 let mut cmd = teaql_core::DeleteCommand::new(key.entity.as_ref(), id);
1293 if let Some(version) = root.get_original_version(key) {
1294 cmd = cmd.expected_version(version);
1295 }
1296 cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1297 self.delete_internal(&cmd).await?;
1298 }
1299
1300 let mut update_batches: std::collections::BTreeMap<
1302 (String, String),
1303 Vec<crate::EntityKey>,
1304 > = std::collections::BTreeMap::new();
1305 let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1306 std::collections::BTreeMap::new();
1307
1308 for (key, record) in &checked_changes {
1309 if deleted_keys.contains(key) {
1310 continue;
1311 }
1312 let mut is_new = new_keys.contains(key);
1313
1314 if !is_new {
1315 let descriptor = self
1316 .data_service
1317 .metadata
1318 .context
1319 .require_entity(&key.entity)
1320 .map_err(DataServiceError::Runtime)?;
1321 let id_property = descriptor.id_property().ok_or_else(|| {
1322 DataServiceError::Runtime(RuntimeError::Graph(format!(
1323 "entity {} has no id property",
1324 key.entity
1325 )))
1326 })?;
1327 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1328 let current_row = self
1329 .fetch_graph_current_row_internal(
1330 &key.entity,
1331 &id_property.name,
1332 &key.id,
1333 my_trace,
1334 )
1335 .await?;
1336 if current_row.is_none() {
1337 is_new = true;
1338 }
1339 }
1340
1341 match is_new {
1342 true => {
1343 insert_batches
1344 .entry(key.entity.to_string())
1345 .or_default()
1346 .push(key.clone());
1347 }
1348 false => {
1349 let mut fields: Vec<String> = record.keys().cloned().collect();
1350 fields.sort();
1351 let signature = fields.join(",");
1352 update_batches
1353 .entry((key.entity.to_string(), signature))
1354 .or_default()
1355 .push(key.clone());
1356 }
1357 }
1358 }
1359
1360 let mut insert_order: Vec<String> = insert_batches.keys().cloned().collect();
1361 insert_order.sort();
1362
1363 for entity in insert_order {
1364 let keys = insert_batches.get(&entity).unwrap();
1365 let descriptor = self
1366 .data_service
1367 .metadata
1368 .context
1369 .require_entity(&entity)
1370 .map_err(DataServiceError::Runtime)?;
1371 let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.name);
1372 let mut traces = Vec::new();
1373 for key in keys {
1374 let record = checked_changes.get(key).unwrap();
1375 let mut db_record = crate::EntityValues::new();
1376 let mut real_id = key.id.clone();
1377 if crate::data_service::helpers::is_unassigned_id_value(&real_id) {
1378 let gen_id = self
1379 .data_service
1380 .metadata
1381 .context
1382 .next_id(&entity)
1383 .map_err(DataServiceError::Runtime)?;
1384 real_id = Value::U64(gen_id);
1385 generated_ids.insert(key.clone(), real_id.clone());
1386 }
1387 db_record.insert("id".to_owned(), real_id);
1388 for (field, value) in record {
1389 if field == "id" {
1390 continue;
1391 }
1392 db_record.insert(field.clone(), value.clone());
1393 }
1394 crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1395 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, true);
1396 cmd.batch_values.push(db_record.into());
1397 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1398 traces.push(my_trace);
1399 }
1400 cmd.trace_chains = traces;
1401 self.execute_prepared_batch_insert(cmd).await?;
1402 }
1403
1404 let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1405 update_order.sort();
1406
1407 for signature in update_order {
1408 let keys = update_batches.get(&signature).unwrap();
1409 let descriptor = self
1410 .data_service
1411 .metadata
1412 .context
1413 .require_entity(&signature.0)
1414 .map_err(DataServiceError::Runtime)?;
1415 let mut update_fields: Vec<String> =
1416 signature.1.split(',').map(|s| s.to_string()).collect();
1417 if descriptor
1418 .properties
1419 .iter()
1420 .any(|p| p.name == "update_time")
1421 && !update_fields.contains(&"update_time".to_owned())
1422 {
1423 update_fields.push("update_time".to_owned());
1424 }
1425 let mut cmd = teaql_core::BatchUpdateCommand::new(&descriptor.name, update_fields);
1426 let mut traces = Vec::new();
1427 for key in keys {
1428 let record = checked_changes.get(key).unwrap();
1429 let mut db_record = crate::EntityValues::new();
1430 db_record.insert("id".to_owned(), key.id.clone());
1431 for (field, value) in record {
1432 if field == "id" {
1433 continue;
1434 }
1435 db_record.insert(field.clone(), value.clone());
1436 }
1437 crate::data_service::helpers::increment_version(
1438 &mut db_record,
1439 descriptor,
1440 root.get_original_version(key),
1441 );
1442 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, false);
1443 cmd.batch_values.push(db_record.into());
1444 cmd.batch_ids.push(key.id.clone());
1445 cmd.batch_expected_versions
1446 .push(root.get_original_version(key));
1447 cmd.batch_old_values.push(None); let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1449 traces.push(my_trace);
1450 }
1451 cmd.trace_chains = traces;
1452 self.execute_prepared_batch_update(cmd).await?;
1453 }
1454
1455 Ok(generated_ids)
1456 }
1457}