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(mut 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 if root.operation != GraphOperation::Remove {
265 let descriptor = self
266 .data_service
267 .metadata
268 .context
269 .require_entity(&root.entity)
270 .map_err(DataServiceError::Runtime)?;
271 let id_property = descriptor.id_property().ok_or_else(|| {
272 DataServiceError::Runtime(RuntimeError::Graph(format!(
273 "entity {} has no id property",
274 root.entity
275 )))
276 })?;
277 let id = root.values.get(&id_property.name).cloned().ok_or_else(|| {
278 DataServiceError::Runtime(RuntimeError::Graph(format!(
279 "saved {} missing identity field {}",
280 root.entity, id_property.name
281 )))
282 })?;
283 root.values = self
284 .fetch_graph_current_row_internal(
285 &root.entity,
286 &id_property.name,
287 &id,
288 root.comment
289 .clone()
290 .map(|comment| {
291 vec![teaql_core::TraceNode {
292 kind: teaql_core::TraceKind::AuditReason,
293 entity_type: root.entity.clone(),
294 entity_id: id.try_u64(),
295 comment,
296 }]
297 })
298 .unwrap_or_default(),
299 )
300 .await?
301 .map(Into::into)
302 .ok_or_else(|| {
303 DataServiceError::Runtime(RuntimeError::Graph(format!(
304 "persisted {} record could not be read back",
305 root.entity
306 )))
307 })?;
308 }
309
310 Ok(root)
311 }
312
313 pub fn graph_node_from_entity<T>(&self, entity: T) -> Result<GraphNode, RuntimeError>
314 where
315 T: Entity,
316 {
317 let descriptor = T::entity_descriptor();
318 if descriptor.name != self.entity {
319 return Err(RuntimeError::Graph(format!(
320 "entity data service {} cannot extract graph root {}",
321 self.entity, descriptor.name
322 )));
323 }
324 let dirty_fields = entity.dirty_fields();
327 let original_values = entity.original_values();
328 let is_deleted = entity.is_marked_as_delete();
329 let comment = entity.get_comment();
330 let mut node = self.graph_node_from_values(&descriptor.name, entity.into_values())?;
331 node.dirty_fields = dirty_fields;
332 node.original_values = original_values.map(Into::into);
333 if is_deleted {
334 node.operation = GraphOperation::Remove;
335 node.relations.clear();
336 }
337 if let Some(c) = comment {
338 node.set_comment(c);
339 }
340 Ok(node)
341 }
342
343 fn collect_graph_plan<'b, 's: 'b>(
344 &'b self,
345 node: &'b mut GraphNode,
346 plan: &'b mut GraphMutationPlan,
347 parent_scope: Option<&'s ScopedCommentNode<'s>>,
348 parent_token: Option<Arc<TraceScopeToken>>,
349 parent_is_create: bool,
350 ) -> std::pin::Pin<
351 Box<dyn std::future::Future<Output = Result<(), DataServiceError<E::Error>>> + Send + '_>,
352 > {
353 Box::pin(async move {
354 match node.operation {
355 GraphOperation::Reference => {
356 plan.push(
357 node.entity.clone(),
358 GraphMutationKind::Reference,
359 node.values.clone().into(),
360 Vec::new(),
361 parent_token,
362 node.original_values.clone(),
363 );
364 return Ok(());
365 }
366 GraphOperation::Remove => {
367 plan.push(
368 node.entity.clone(),
369 GraphMutationKind::Delete,
370 node.values.clone().into(),
371 Vec::new(),
372 parent_token,
373 node.original_values.clone(),
374 );
375 return Ok(());
376 }
377 GraphOperation::Upsert | GraphOperation::Create => {}
378 }
379
380 let descriptor = self
381 .data_service
382 .metadata
383 .context
384 .require_entity(&node.entity)
385 .map_err(DataServiceError::Runtime)?;
386
387 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
389 parent: parent_scope,
390 track: teaql_core::TraceNode {
391 kind: teaql_core::TraceKind::AuditReason,
392 entity_type: node.entity.clone(),
393 entity_id: node.id().and_then(|v| match v {
394 Value::U64(n) => Some(*n),
395 Value::I64(n) => Some(*n as u64),
396 _ => None,
397 }),
398 comment: c.clone(),
399 },
400 });
401 let active_scope = current_scope.as_ref().or(parent_scope);
402
403 let id_property = descriptor.id_property().cloned();
404 let id = id_property.as_ref().and_then(|property| {
405 node.values
406 .get(&property.name)
407 .filter(|value| !is_unassigned_id_value(value))
408 .cloned()
409 });
410
411 if let Some(id_val) = &id {
412 if !plan
413 .visited_nodes
414 .insert((node.entity.clone(), graph_identity_key(id_val)))
415 {
416 return Ok(());
417 }
418 }
419
420 let is_create_op = node.operation == GraphOperation::Create
421 || (parent_is_create && node.operation == GraphOperation::Upsert);
422
423 let is_update = match is_create_op {
424 true => false,
425 false => match (id_property.as_ref(), id.as_ref()) {
426 (Some(id_property), Some(id)) => self
427 .fetch_graph_current_row_internal(
428 &node.entity,
429 &id_property.name,
430 id,
431 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
432 )
433 .await?
434 .is_some(),
435 _ => false,
436 },
437 };
438 if !is_update {
439 if let Some(id_property) = id_property.as_ref() {
440 let needs_id = !node.values.contains_key(&id_property.name)
441 || node
442 .values
443 .get(&id_property.name)
444 .is_some_and(is_unassigned_id_value);
445 if needs_id {
446 let id = self
447 .data_service
448 .metadata
449 .context
450 .next_id(&node.entity)
451 .map_err(DataServiceError::Runtime)?;
452 node.values.insert(id_property.name.clone(), Value::U64(id));
453 }
454 }
455 ensure_initial_version(&mut node.values, descriptor);
456 crate::data_service::helpers::ensure_timestamps(&mut node.values, descriptor, true);
457 } else {
458 crate::data_service::helpers::ensure_timestamps(
459 &mut node.values,
460 descriptor,
461 false,
462 );
463 }
464 let update_fields = is_update
465 .then(|| {
466 let mut excluded = Vec::new();
467 if let Some(id_property) = id_property.as_ref() {
468 excluded.push(id_property.name.clone());
469 }
470 if let Some(version_property) = descriptor.version_property() {
471 excluded.push(version_property.name.clone());
472 }
473 let mut fields = sorted_update_fields(&node.values, excluded);
474 if let Some(dirty) = &node.dirty_fields {
475 fields.retain(|f| dirty.contains(f));
476 }
477 fields
478 })
479 .unwrap_or_default();
480
481 let current_token = node
484 .comment
485 .as_ref()
486 .map(|c| {
487 Arc::new(TraceScopeToken {
488 parent: parent_token.clone(),
489 track: teaql_core::TraceNode {
490 kind: teaql_core::TraceKind::AuditReason,
491 entity_type: node.entity.clone(),
492 entity_id: node.id().and_then(|v| match v {
493 Value::U64(n) => Some(*n),
494 Value::I64(n) => Some(*n as u64),
495 _ => None,
496 }),
497 comment: c.clone(),
498 },
499 node_index: plan.next_item_index,
500 })
501 })
502 .or_else(|| parent_token.clone());
503
504 plan.push(
505 node.entity.clone(),
506 GraphMutationKind::for_update(is_update),
507 node.values.clone().into(),
508 update_fields,
509 current_token.clone(),
510 node.original_values.clone(),
511 );
512
513 for (name, children) in &mut node.relations {
514 let relation = descriptor.relation_by_name(name).ok_or_else(|| {
515 DataServiceError::Runtime(RuntimeError::MissingRelation {
516 entity: node.entity.clone(),
517 relation: name.clone(),
518 })
519 })?;
520 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
521 for child in children {
522 ensure_relation_target(&node.entity, name, &relation.target_entity, child)?;
523 child_repo
524 .collect_graph_plan(
525 child,
526 plan,
527 active_scope,
528 current_token.clone(),
529 is_create_op,
530 )
531 .await?;
532 }
533 }
534 Ok(())
535 })
536 }
537
538 fn insert_graph_node_scoped<'b, 's: 'b>(
539 &'b self,
540 mut node: GraphNode,
541 parent_scope: Option<&'s ScopedCommentNode<'s>>,
542 ) -> std::pin::Pin<
543 Box<
544 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
545 + Send
546 + '_,
547 >,
548 > {
549 Box::pin(async move {
550 match node.operation {
551 GraphOperation::Upsert | GraphOperation::Create => {}
552 GraphOperation::Reference => {
553 return self
554 .validate_reference_node(
555 node,
556 parent_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
557 )
558 .await;
559 }
560 GraphOperation::Remove => {
561 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
562 "create graph cannot remove node {}",
563 node.entity
564 ))));
565 }
566 }
567
568 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
570 parent: parent_scope,
571 track: teaql_core::TraceNode {
572 kind: teaql_core::TraceKind::AuditReason,
573 entity_type: node.entity.clone(),
574 entity_id: node.id().and_then(|v| match v {
575 Value::U64(n) => Some(*n),
576 Value::I64(n) => Some(*n as u64),
577 _ => None,
578 }),
579 comment: c.clone(),
580 },
581 });
582 let active_scope = current_scope.as_ref().or(parent_scope);
583
584 let descriptor = self
585 .data_service
586 .metadata
587 .context
588 .require_entity(&node.entity)
589 .map_err(DataServiceError::Runtime)?;
590
591 let mut one_relations = Vec::new();
592 let mut many_relations = Vec::new();
593 for (name, children) in std::mem::take(&mut node.relations) {
594 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
595 DataServiceError::Runtime(RuntimeError::MissingRelation {
596 entity: node.entity.clone(),
597 relation: name.clone(),
598 })
599 })?;
600 match relation.many {
601 true => many_relations.push((name, relation.clone(), children)),
602 false => one_relations.push((name, relation.clone(), children)),
603 }
604 }
605
606 for (name, relation, children) in one_relations {
607 if children.len() > 1 {
608 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
609 "relation {}.{} expects one child, got {}",
610 node.entity,
611 name,
612 children.len()
613 ))));
614 }
615 let mut saved_children = Vec::new();
616 for child in children {
617 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
618 let child_repo = self.scoped_data_service_internal(child.entity.clone());
619 let saved_child = child_repo
620 .insert_graph_node_scoped(child, active_scope)
621 .await?;
622 if relation.attach {
623 let foreign_value = saved_child
624 .values
625 .get(&relation.foreign_key)
626 .cloned()
627 .ok_or_else(|| {
628 DataServiceError::Runtime(RuntimeError::Graph(format!(
629 "saved child {} missing foreign key {} for relation {}.{}",
630 relation.target_entity, relation.foreign_key, node.entity, name
631 )))
632 })?;
633 node.values
634 .insert(relation.local_key.clone(), foreign_value);
635 }
636 saved_children.push(saved_child);
637 }
638 node.relations.insert(name, saved_children);
639 }
640
641 let command = self
642 .prepare_insert_command(&InsertCommand {
643 entity: node.entity.clone(),
644 values: node.values.clone().into(),
645 trace_chain: Vec::new(),
646 })
647 .map_err(DataServiceError::Runtime)?;
648 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
649 self.execute_prepared_insert_with_comment(command.clone(), lineage)
650 .await?;
651 node.values = command.values.into();
652 if let Some(id_property) = descriptor.id_property() {
653 if let Some(id) = node.values.get(&id_property.name).cloned() {
654 node.values = self
655 .fetch_graph_current_row_internal(
656 &node.entity,
657 &id_property.name,
658 &id,
659 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
660 )
661 .await?
662 .map(Into::into)
663 .ok_or_else(|| {
664 DataServiceError::Runtime(RuntimeError::Graph(format!(
665 "persisted {} record could not be read back",
666 node.entity
667 )))
668 })?;
669 }
670 }
671
672 for (name, relation, children) in many_relations {
673 let local_value =
674 node.values
675 .get(&relation.local_key)
676 .cloned()
677 .ok_or_else(|| {
678 DataServiceError::Runtime(RuntimeError::Graph(format!(
679 "parent {} missing local key {} for relation {}",
680 node.entity, relation.local_key, name
681 )))
682 })?;
683 let mut saved_children = Vec::new();
684 for mut child in children {
685 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
686 if relation.attach {
687 child
688 .values
689 .insert(relation.foreign_key.clone(), local_value.clone());
690 }
691 let child_repo = self.scoped_data_service_internal(child.entity.clone());
692 saved_children.push(
693 child_repo
694 .insert_graph_node_scoped(child, active_scope)
695 .await?,
696 );
697 }
698 node.relations.insert(name, saved_children);
699 }
700
701 Ok(node)
702 })
703 }
704
705 fn upsert_graph_node_scoped<'b, 's: 'b>(
706 &'b self,
707 mut node: GraphNode,
708 parent_scope: Option<&'s ScopedCommentNode<'s>>,
709 ) -> std::pin::Pin<
710 Box<
711 dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
712 + Send
713 + '_,
714 >,
715 > {
716 Box::pin(async move {
717 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
719 parent: parent_scope,
720 track: teaql_core::TraceNode {
721 kind: teaql_core::TraceKind::AuditReason,
722 entity_type: node.entity.clone(),
723 entity_id: node.id().and_then(|v| match v {
724 Value::U64(n) => Some(*n),
725 Value::I64(n) => Some(*n as u64),
726 _ => None,
727 }),
728 comment: c.clone(),
729 },
730 });
731 let active_scope = current_scope.as_ref().or(parent_scope);
732
733 match node.operation {
734 GraphOperation::Upsert | GraphOperation::Create => {}
735 GraphOperation::Reference => {
736 return self
737 .validate_reference_node(
738 node,
739 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
740 )
741 .await;
742 }
743 GraphOperation::Remove => {
744 self.validate_remove_node(
745 &node,
746 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
747 )
748 .await?;
749 self.delete_graph_node(&node, parent_scope).await?;
750 return Ok(node);
751 }
752 }
753
754 let descriptor = self
755 .data_service
756 .metadata
757 .context
758 .require_entity(&node.entity)
759 .map_err(DataServiceError::Runtime)?;
760 let Some(id_property) = descriptor.id_property() else {
761 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
762 "entity {} has no id property for graph upsert",
763 node.entity
764 ))));
765 };
766 let Some(id) = node
767 .values
768 .get(&id_property.name)
769 .filter(|value| !is_unassigned_id_value(value))
770 .cloned()
771 else {
772 node.comment = None;
774 return self.insert_graph_node_scoped(node, active_scope).await;
775 };
776
777 if node.operation == GraphOperation::Create
778 || self
779 .fetch_graph_current_row_internal(
780 &node.entity,
781 &id_property.name,
782 &id,
783 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
784 )
785 .await?
786 .is_none()
787 {
788 node.comment = None;
789 return self.insert_graph_node_scoped(node, active_scope).await;
790 }
791
792 let mut one_relations = Vec::new();
793 let mut many_relations = Vec::new();
794 for (name, children) in std::mem::take(&mut node.relations) {
795 let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
796 DataServiceError::Runtime(RuntimeError::MissingRelation {
797 entity: node.entity.clone(),
798 relation: name.clone(),
799 })
800 })?;
801 match relation.many {
802 true => many_relations.push((name, relation.clone(), children)),
803 false => one_relations.push((name, relation.clone(), children)),
804 }
805 }
806
807 for (name, relation, children) in one_relations {
808 if children.len() > 1 {
809 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
810 "relation {}.{} expects one child, got {}",
811 node.entity,
812 name,
813 children.len()
814 ))));
815 }
816 let mut saved_children = Vec::new();
817 for child in children {
818 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
819 let child_repo = self.scoped_data_service_internal(child.entity.clone());
820 let saved_child = child_repo
821 .upsert_graph_node_scoped(child, active_scope)
822 .await?;
823 if relation.attach {
824 let foreign_value = saved_child
825 .values
826 .get(&relation.foreign_key)
827 .cloned()
828 .ok_or_else(|| {
829 DataServiceError::Runtime(RuntimeError::Graph(format!(
830 "saved child {} missing foreign key {} for relation {}.{}",
831 relation.target_entity, relation.foreign_key, node.entity, name
832 )))
833 })?;
834 node.values
835 .insert(relation.local_key.clone(), foreign_value);
836 }
837 saved_children.push(saved_child);
838 }
839 node.relations.insert(name, saved_children);
840 }
841
842 let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
843 if !update.values.is_empty() {
844 let prepared_update = self
845 .prepare_update_command(&update)
846 .map_err(DataServiceError::Runtime)?;
847 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
848 self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
849 .await?;
850 for (field, value) in &prepared_update.values {
851 node.values.insert(field.clone(), value.clone());
852 }
853 if let Some(version_property) = descriptor.version_property() {
854 if let Some(expected_version) = prepared_update.expected_version {
855 node.values.insert(
856 version_property.name.clone(),
857 Value::I64(expected_version + 1),
858 );
859 }
860 }
861 node.values = self
862 .fetch_graph_current_row_internal(
863 &node.entity,
864 &id_property.name,
865 &id,
866 active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
867 )
868 .await?
869 .map(Into::into)
870 .ok_or_else(|| {
871 DataServiceError::Runtime(RuntimeError::Graph(format!(
872 "persisted {} record could not be read back",
873 node.entity
874 )))
875 })?;
876 }
877
878 for (name, relation, children) in many_relations {
879 let local_value =
880 node.values
881 .get(&relation.local_key)
882 .cloned()
883 .ok_or_else(|| {
884 DataServiceError::Runtime(RuntimeError::Graph(format!(
885 "parent {} missing local key {} for relation {}",
886 node.entity, relation.local_key, name
887 )))
888 })?;
889 let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
890 let child_descriptor = self
891 .data_service
892 .metadata
893 .context
894 .require_entity(&relation.target_entity)
895 .map_err(DataServiceError::Runtime)?;
896 let child_id_property = child_descriptor.id_property().ok_or_else(|| {
897 DataServiceError::Runtime(RuntimeError::Graph(format!(
898 "entity {} has no id property",
899 relation.target_entity
900 )))
901 })?;
902
903 let mut seen = std::collections::BTreeSet::new();
904 let mut saved_children = Vec::new();
905 for mut child in children {
906 ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
907 if relation.attach && child.operation != GraphOperation::Reference {
908 child
909 .values
910 .insert(relation.foreign_key.clone(), local_value.clone());
911 }
912 if let Some(child_id) = child
913 .values
914 .get(&child_id_property.name)
915 .filter(|value| !is_unassigned_id_value(value))
916 {
917 let key = graph_identity_key(child_id);
918 if !seen.insert(key.clone()) {
919 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
920 "duplicate child id {key} in relation {}.{}",
921 node.entity, name
922 ))));
923 }
924 }
925 saved_children.push(
926 child_repo
927 .upsert_graph_node_scoped(child, active_scope)
928 .await?,
929 );
930 }
931
932 node.relations.insert(name, saved_children);
933 }
934
935 Ok(node)
936 })
937 }
938
939 async fn validate_reference_node(
940 &self,
941 node: GraphNode,
942 trace_chain: Vec<teaql_core::TraceNode>,
943 ) -> Result<GraphNode, DataServiceError<E::Error>> {
944 if !node.relations.is_empty() {
945 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
946 "reference node {} cannot contain child relations",
947 node.entity
948 ))));
949 }
950 let descriptor = self
951 .data_service
952 .metadata
953 .context
954 .require_entity(&node.entity)
955 .map_err(DataServiceError::Runtime)?;
956 let id_property = descriptor.id_property().ok_or_else(|| {
957 DataServiceError::Runtime(RuntimeError::Graph(format!(
958 "entity {} has no id property for graph reference",
959 node.entity
960 )))
961 })?;
962 let id = node
963 .values
964 .get(&id_property.name)
965 .filter(|value| !is_unassigned_id_value(value))
966 .cloned()
967 .ok_or_else(|| {
968 DataServiceError::Runtime(RuntimeError::Graph(format!(
969 "reference node {} missing id property {}",
970 node.entity, id_property.name
971 )))
972 })?;
973
974 for field in node.values.keys() {
975 if field == &id_property.name {
976 continue;
977 }
978 if descriptor
979 .version_property()
980 .map(|property| field == &property.name)
981 .unwrap_or(false)
982 {
983 continue;
984 }
985 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
986 "reference node {} cannot carry mutable field {}",
987 node.entity, field
988 ))));
989 }
990
991 let current = self
992 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
993 .await?
994 .ok_or_else(|| {
995 DataServiceError::Runtime(RuntimeError::Graph(format!(
996 "reference node {}({}) does not exist",
997 node.entity,
998 graph_identity_key(&id)
999 )))
1000 })?;
1001
1002 if let Some(version_property) = descriptor.version_property() {
1003 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
1004 if *existing_version < 0 {
1005 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1006 "reference node {}({}) is deleted",
1007 node.entity,
1008 graph_identity_key(&id)
1009 ))));
1010 }
1011 if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
1012 {
1013 if expected_version != existing_version {
1014 println!(
1015 "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
1016 node.entity, expected_version, existing_version
1017 );
1018 return Err(DataServiceError::Runtime(
1019 RuntimeError::OptimisticLockConflict {
1020 entity: node.entity,
1021 id: graph_identity_key(&id),
1022 },
1023 ));
1024 }
1025 }
1026 }
1027 }
1028
1029 Ok(GraphNode {
1030 entity: node.entity,
1031 values: current.into(),
1032 relations: BTreeMap::new(),
1033 operation: GraphOperation::Reference,
1034 comment: None,
1035 dirty_fields: None,
1036 original_values: None,
1037 })
1038 }
1039
1040 async fn validate_remove_node(
1041 &self,
1042 node: &GraphNode,
1043 trace_chain: Vec<teaql_core::TraceNode>,
1044 ) -> Result<(), DataServiceError<E::Error>> {
1045 if !node.relations.is_empty() {
1046 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1047 "remove node {} cannot contain child relations",
1048 node.entity
1049 ))));
1050 }
1051 let descriptor = self
1052 .data_service
1053 .metadata
1054 .context
1055 .require_entity(&node.entity)
1056 .map_err(DataServiceError::Runtime)?;
1057 let id_property = descriptor.id_property().ok_or_else(|| {
1058 DataServiceError::Runtime(RuntimeError::Graph(format!(
1059 "entity {} has no id property for graph remove",
1060 node.entity
1061 )))
1062 })?;
1063 let id = node
1064 .values
1065 .get(&id_property.name)
1066 .filter(|value| !is_unassigned_id_value(value))
1067 .cloned()
1068 .ok_or_else(|| {
1069 DataServiceError::Runtime(RuntimeError::Graph(format!(
1070 "remove node {} missing id property {}",
1071 node.entity, id_property.name
1072 )))
1073 })?;
1074 let current = self
1075 .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
1076 .await?
1077 .ok_or_else(|| {
1078 DataServiceError::Runtime(RuntimeError::Graph(format!(
1079 "remove node {}({}) does not exist",
1080 node.entity,
1081 graph_identity_key(&id)
1082 )))
1083 })?;
1084 if let Some(version_property) = descriptor.version_property() {
1085 if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
1086 if *existing_version < 0 {
1087 return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1088 "remove node {}({}) is already deleted",
1089 node.entity,
1090 graph_identity_key(&id)
1091 ))));
1092 }
1093 }
1094 }
1095 Ok(())
1096 }
1097
1098 fn graph_node_from_values(
1099 &self,
1100 entity: &str,
1101 values: MutationValues,
1102 ) -> Result<GraphNode, RuntimeError> {
1103 let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1104 let mut node = GraphNode::new(entity);
1105
1106 for (field, value) in values {
1107 if field == "_comment" {
1108 if let Value::Text(comment) = value {
1109 node.set_comment(comment);
1110 }
1111 continue;
1112 }
1113 if field == "_dirty_fields" {
1114 if let Value::List(fields) = value {
1115 let mut dirty = std::collections::BTreeSet::new();
1116 for f in fields {
1117 if let Value::Text(t) = f {
1118 dirty.insert(t);
1119 }
1120 }
1121 node.dirty_fields = Some(dirty);
1122 }
1123 continue;
1124 }
1125 if field == "_original_values" {
1126 if let Value::Object(orig) = value {
1127 node.original_values = Some(orig.into());
1128 }
1129 continue;
1130 }
1131 if field == "_is_new" {
1132 if matches!(value, Value::Bool(true)) {
1133 node.operation = GraphOperation::Create;
1134 }
1135 continue;
1136 }
1137 if field == "_is_deleted" {
1138 if matches!(value, Value::Bool(true)) {
1139 node.operation = GraphOperation::Remove;
1140 }
1141 continue;
1142 }
1143 let Some(relation) = descriptor.relation_by_name(&field) else {
1144 node.values.insert(field, value);
1145 continue;
1146 };
1147
1148 match value {
1149 Value::Null => {
1150 node.relations.entry(field).or_default();
1151 }
1152 Value::Object(record) => {
1153 let child =
1154 self.graph_node_from_values(&relation.target_entity, record.into())?;
1155 node.relations.entry(field).or_default().push(child);
1156 }
1157 Value::List(values) => {
1158 let children = node.relations.entry(field.clone()).or_default();
1159 for value in values {
1160 let Value::Object(record) = value else {
1161 return Err(RuntimeError::Graph(format!(
1162 "relation {}.{} expects object children, got {:?}",
1163 entity, field, value
1164 )));
1165 };
1166 children.push(
1167 self.graph_node_from_values(&relation.target_entity, record.into())?,
1168 );
1169 }
1170 }
1171 other => {
1172 return Err(RuntimeError::Graph(format!(
1173 "relation {}.{} expects object/list/null, got {:?}",
1174 entity, field, other
1175 )));
1176 }
1177 }
1178 }
1179
1180 Ok(node)
1181 }
1182
1183 fn graph_update_command(
1184 &self,
1185 node: &mut GraphNode,
1186 descriptor: &EntityDescriptor,
1187 id_property: &PropertyDescriptor,
1188 id: &Value,
1189 ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1190 crate::mark_entity_status(&mut node.values, crate::CheckObjectStatus::Update);
1191 let check_result = self
1192 .data_service
1193 .metadata
1194 .context
1195 .check_and_fix_values(&node.entity, &mut node.values);
1196 crate::clear_entity_status(&mut node.values);
1197 check_result.map_err(DataServiceError::Runtime)?;
1198
1199 let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1200 command.old_values = node.original_values.clone().map(Into::into);
1201 if let Some(version_property) = descriptor.version_property() {
1202 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1203 command = command.expected_version(*version);
1204 }
1205 }
1206 for property in descriptor.properties.iter().filter(|property| {
1210 !property.is_id
1211 && !property.is_version
1212 && property.name != id_property.name
1213 && match &node.dirty_fields {
1214 Some(dirty) => dirty.contains(&property.name),
1215 None => node.values.contains_key(&property.name),
1216 }
1217 }) {
1218 if let Some(value) = node.values.get(&property.name) {
1219 command.values.insert(property.name.clone(), value.clone());
1220 }
1221 }
1222 Ok(command)
1223 }
1224
1225 fn delete_graph_node<'b, 's: 'b>(
1226 &'b self,
1227 node: &'b GraphNode,
1228 parent_scope: Option<&'s ScopedCommentNode<'s>>,
1229 ) -> std::pin::Pin<
1230 Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1231 > {
1232 Box::pin(async move {
1233 let descriptor = self
1234 .data_service
1235 .metadata
1236 .context
1237 .require_entity(&node.entity)
1238 .map_err(DataServiceError::Runtime)?;
1239 let id_property = descriptor.id_property().ok_or_else(|| {
1240 DataServiceError::Runtime(RuntimeError::Graph(format!(
1241 "entity {} has no id property for graph remove",
1242 node.entity
1243 )))
1244 })?;
1245 let id = node
1246 .values
1247 .get(&id_property.name)
1248 .filter(|value| !is_unassigned_id_value(value))
1249 .cloned()
1250 .ok_or_else(|| {
1251 DataServiceError::Runtime(RuntimeError::Graph(format!(
1252 "remove node {} missing id property {}",
1253 node.entity, id_property.name
1254 )))
1255 })?;
1256 let mut delete = DeleteCommand::new(node.entity.clone(), id);
1257 if let Some(version_property) = descriptor.version_property() {
1258 if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1259 delete = delete.expected_version(*version);
1260 }
1261 }
1262
1263 let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1265 parent: parent_scope,
1266 track: teaql_core::TraceNode {
1267 kind: teaql_core::TraceKind::AuditReason,
1268 entity_type: node.entity.clone(),
1269 entity_id: node.id().and_then(|v| match v {
1270 Value::U64(n) => Some(*n),
1271 Value::I64(n) => Some(*n as u64),
1272 _ => None,
1273 }),
1274 comment: c.clone(),
1275 },
1276 });
1277 let active_scope = current_scope.as_ref().or(parent_scope);
1278 let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1279
1280 self.delete_scoped_internal(&delete, lineage).await
1281 })
1282 }
1283
1284 pub(crate) async fn fetch_graph_current_row_internal(
1285 &self,
1286 entity: &str,
1287 id_property: &str,
1288 id: &teaql_core::Value,
1289 trace_chain: Vec<teaql_core::TraceNode>,
1290 ) -> Result<Option<teaql_core::CompactRow>, DataServiceError<E::Error>> {
1291 let mut query = teaql_core::SelectQuery::new(entity)
1292 .filter(teaql_core::Expr::eq(id_property, id.clone()));
1293 query.trace_chain = trace_chain;
1294 let mut rows = self
1295 .scoped_data_service_internal(entity.to_owned())
1296 .fetch_all_internal(&query)
1297 .await?;
1298 Ok(rows.pop())
1299 }
1300
1301 pub(crate) async fn execute_ledger_plan_internal(
1302 &self,
1303 root: crate::EntityRuntimeState,
1304 ) -> Result<std::collections::BTreeMap<crate::EntityKey, Value>, DataServiceError<E::Error>>
1305 {
1306 let mut generated_ids = std::collections::BTreeMap::new();
1307 let comment = root.get_comment();
1308 let trace_chain = comment
1309 .map(|c| {
1310 vec![teaql_core::TraceNode {
1311 kind: teaql_core::TraceKind::AuditReason,
1312 entity_type: self.entity.clone(),
1313 entity_id: None,
1314 comment: c,
1315 }]
1316 })
1317 .unwrap_or_default();
1318
1319 let deleted_keys = root.deleted_keys();
1320 let new_keys = root.new_keys();
1321 let change_set = root.current_change_set();
1322
1323 let mut checked_changes = std::collections::BTreeMap::new();
1330 for (key, record) in change_set.changes() {
1331 if deleted_keys.contains(key) {
1332 continue;
1333 }
1334 let mut checked: crate::EntityValues = record.clone().into();
1335 checked
1336 .entry("id".to_owned())
1337 .or_insert_with(|| key.id.clone());
1338 checked_changes.insert(key.clone(), checked);
1339 }
1340
1341 for key in deleted_keys.iter() {
1343 let id = key.id.clone();
1344 let mut cmd = teaql_core::DeleteCommand::new(key.entity.as_ref(), id);
1345 if let Some(version) = root.get_original_version(key) {
1346 cmd = cmd.expected_version(version);
1347 }
1348 cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1349 self.delete_internal(&cmd).await?;
1350 }
1351
1352 let mut update_batches: std::collections::BTreeMap<
1354 (String, String),
1355 Vec<crate::EntityKey>,
1356 > = std::collections::BTreeMap::new();
1357 let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1358 std::collections::BTreeMap::new();
1359
1360 for (key, record) in &checked_changes {
1361 if deleted_keys.contains(key) {
1362 continue;
1363 }
1364 let mut is_new = new_keys.contains(key);
1365
1366 if !is_new {
1367 let descriptor = self
1368 .data_service
1369 .metadata
1370 .context
1371 .require_entity(&key.entity)
1372 .map_err(DataServiceError::Runtime)?;
1373 let id_property = descriptor.id_property().ok_or_else(|| {
1374 DataServiceError::Runtime(RuntimeError::Graph(format!(
1375 "entity {} has no id property",
1376 key.entity
1377 )))
1378 })?;
1379 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1380 let current_row = self
1381 .fetch_graph_current_row_internal(
1382 &key.entity,
1383 &id_property.name,
1384 &key.id,
1385 my_trace,
1386 )
1387 .await?;
1388 if current_row.is_none() {
1389 is_new = true;
1390 }
1391 }
1392
1393 match is_new {
1394 true => {
1395 insert_batches
1396 .entry(key.entity.to_string())
1397 .or_default()
1398 .push(key.clone());
1399 }
1400 false => {
1401 let mut fields: Vec<String> = record.keys().cloned().collect();
1402 fields.sort();
1403 let signature = fields.join(",");
1404 update_batches
1405 .entry((key.entity.to_string(), signature))
1406 .or_default()
1407 .push(key.clone());
1408 }
1409 }
1410 }
1411
1412 let mut insert_order: Vec<String> = insert_batches.keys().cloned().collect();
1413 insert_order.sort();
1414
1415 for entity in insert_order {
1416 let keys = insert_batches.get(&entity).unwrap();
1417 let descriptor = self
1418 .data_service
1419 .metadata
1420 .context
1421 .require_entity(&entity)
1422 .map_err(DataServiceError::Runtime)?;
1423 let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.name);
1424 let mut traces = Vec::new();
1425 for key in keys {
1426 let record = checked_changes.get(key).unwrap();
1427 let mut db_record = crate::EntityValues::new();
1428 let mut real_id = key.id.clone();
1429 if crate::data_service::helpers::is_unassigned_id_value(&real_id) {
1430 let gen_id = self
1431 .data_service
1432 .metadata
1433 .context
1434 .next_id(&entity)
1435 .map_err(DataServiceError::Runtime)?;
1436 real_id = Value::U64(gen_id);
1437 generated_ids.insert(key.clone(), real_id.clone());
1438 }
1439 db_record.insert("id".to_owned(), real_id);
1440 for (field, value) in record {
1441 if field == "id" {
1442 continue;
1443 }
1444 db_record.insert(field.clone(), value.clone());
1445 }
1446 crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1447 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, true);
1448 cmd.batch_values.push(db_record.into());
1449 let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1450 traces.push(my_trace);
1451 }
1452 cmd.trace_chains = traces;
1453 self.execute_prepared_batch_insert(cmd).await?;
1454 }
1455
1456 let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1457 update_order.sort();
1458
1459 for signature in update_order {
1460 let keys = update_batches.get(&signature).unwrap();
1461 let descriptor = self
1462 .data_service
1463 .metadata
1464 .context
1465 .require_entity(&signature.0)
1466 .map_err(DataServiceError::Runtime)?;
1467 let mut update_fields: Vec<String> =
1468 signature.1.split(',').map(|s| s.to_string()).collect();
1469 if descriptor
1470 .properties
1471 .iter()
1472 .any(|p| p.name == "update_time")
1473 && !update_fields.contains(&"update_time".to_owned())
1474 {
1475 update_fields.push("update_time".to_owned());
1476 }
1477 let mut cmd = teaql_core::BatchUpdateCommand::new(&descriptor.name, update_fields);
1478 let mut traces = Vec::new();
1479 for key in keys {
1480 let record = checked_changes.get(key).unwrap();
1481 let mut db_record = crate::EntityValues::new();
1482 db_record.insert("id".to_owned(), key.id.clone());
1483 for (field, value) in record {
1484 if field == "id" {
1485 continue;
1486 }
1487 db_record.insert(field.clone(), value.clone());
1488 }
1489 crate::data_service::helpers::increment_version(
1490 &mut db_record,
1491 descriptor,
1492 root.get_original_version(key),
1493 );
1494 crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, false);
1495 cmd.batch_values.push(db_record.into());
1496 cmd.batch_ids.push(key.id.clone());
1497 cmd.batch_expected_versions
1498 .push(root.get_original_version(key));
1499 cmd.batch_old_values.push(None); let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1501 traces.push(my_trace);
1502 }
1503 cmd.trace_chains = traces;
1504 self.execute_prepared_batch_update(cmd).await?;
1505 }
1506
1507 Ok(generated_ids)
1508 }
1509}