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