Skip to main content

teaql_runtime/data_service/
graph.rs

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