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
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    /// Create a new entity graph with an annotation comment on the root node.
155    /// This assumes all new nodes do not exist in the database, skipping existence checks
156    /// and throwing an exception on primary key conflict.
157    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 now, loop individually since we lack BatchDeleteCommand
247                    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                    // References are skipped in execution, they only validate during traversal
264                }
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        // Extract dirty field names BEFORE into_record() consumes the entity.
283        // This is the Rust equivalent of Java's entity.getUpdatedProperties().
284        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            // Create scope node on the current stack frame if this node has a comment
346            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            }
414            let update_fields = is_update
415                .then(|| {
416                    let mut excluded = Vec::new();
417                    if let Some(id_property) = id_property.as_ref() {
418                        excluded.push(id_property.name.clone());
419                    }
420                    if let Some(version_property) = descriptor.version_property() {
421                        excluded.push(version_property.name.clone());
422                    }
423                    let mut fields = sorted_update_fields(&node.values, excluded);
424                    if let Some(dirty) = &node.dirty_fields {
425                        fields.retain(|f| dirty.contains(f));
426                    }
427                    fields
428                })
429                .unwrap_or_default();
430
431            // Build the TraceScopeToken for this node (only if it has a comment).
432            // This is an Arc-linked persistent list: zero-copy, O(1) creation.
433            let current_token = node
434                .comment
435                .as_ref()
436                .map(|c| {
437                    Arc::new(TraceScopeToken {
438                        parent: parent_token.clone(),
439                        track: teaql_core::TraceNode {
440                            entity_type: node.entity.clone(),
441                            entity_id: node.id().and_then(|v| match v {
442                                Value::U64(n) => Some(*n),
443                                Value::I64(n) => Some(*n as u64),
444                                _ => None,
445                            }),
446                            comment: c.clone(),
447                        },
448                        node_index: plan.next_item_index,
449                    })
450                })
451                .or_else(|| parent_token.clone());
452
453            plan.push(
454                node.entity.clone(),
455                GraphMutationKind::for_update(is_update),
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_internal(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                match relation.many {
549                    true => many_relations.push((name, relation.clone(), children)),
550                    false => one_relations.push((name, relation.clone(), children)),
551                }
552            }
553
554            for (name, relation, children) in one_relations {
555                if children.len() > 1 {
556                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
557                        "relation {}.{} expects one child, got {}",
558                        node.entity,
559                        name,
560                        children.len()
561                    ))));
562                }
563                let mut saved_children = Vec::new();
564                for child in children {
565                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
566                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
567                    let saved_child = child_repo
568                        .insert_graph_node_scoped(child, active_scope)
569                        .await?;
570                    if relation.attach {
571                        let foreign_value = saved_child
572                            .values
573                            .get(&relation.foreign_key)
574                            .cloned()
575                            .ok_or_else(|| {
576                                DataServiceError::Runtime(RuntimeError::Graph(format!(
577                                    "saved child {} missing foreign key {} for relation {}.{}",
578                                    relation.target_entity, relation.foreign_key, node.entity, name
579                                )))
580                            })?;
581                        node.values
582                            .insert(relation.local_key.clone(), foreign_value);
583                    }
584                    saved_children.push(saved_child);
585                }
586                node.relations.insert(name, saved_children);
587            }
588
589            let command = self
590                .prepare_insert_command(&InsertCommand {
591                    entity: node.entity.clone(),
592                    values: node.values.clone(),
593                    trace_chain: Vec::new(),
594                })
595                .map_err(DataServiceError::Runtime)?;
596            let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
597            self.execute_prepared_insert_with_comment(command.clone(), lineage)
598                .await?;
599            node.values = command.values;
600
601            for (name, relation, children) in many_relations {
602                let local_value =
603                    node.values
604                        .get(&relation.local_key)
605                        .cloned()
606                        .ok_or_else(|| {
607                            DataServiceError::Runtime(RuntimeError::Graph(format!(
608                                "parent {} missing local key {} for relation {}",
609                                node.entity, relation.local_key, name
610                            )))
611                        })?;
612                let mut saved_children = Vec::new();
613                for mut child in children {
614                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
615                    if relation.attach {
616                        child
617                            .values
618                            .insert(relation.foreign_key.clone(), local_value.clone());
619                    }
620                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
621                    saved_children.push(
622                        child_repo
623                            .insert_graph_node_scoped(child, active_scope)
624                            .await?,
625                    );
626                }
627                node.relations.insert(name, saved_children);
628            }
629
630            Ok(node)
631        })
632    }
633
634    fn upsert_graph_node_scoped<'b, 's: 'b>(
635        &'b self,
636        mut node: GraphNode,
637        parent_scope: Option<&'s ScopedCommentNode<'s>>,
638    ) -> std::pin::Pin<
639        Box<
640            dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
641                + Send
642                + '_,
643        >,
644    > {
645        Box::pin(async move {
646            // Create scope node on the current stack frame if this node has a comment
647            let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
648                parent: parent_scope,
649                track: teaql_core::TraceNode {
650                    entity_type: node.entity.clone(),
651                    entity_id: node.id().and_then(|v| match v {
652                        Value::U64(n) => Some(*n),
653                        Value::I64(n) => Some(*n as u64),
654                        _ => None,
655                    }),
656                    comment: c.clone(),
657                },
658            });
659            let active_scope = current_scope.as_ref().or(parent_scope);
660
661            match node.operation {
662                GraphOperation::Upsert | GraphOperation::Create => {}
663                GraphOperation::Reference => {
664                    return self
665                        .validate_reference_node(
666                            node,
667                            active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
668                        )
669                        .await;
670                }
671                GraphOperation::Remove => {
672                    self.validate_remove_node(
673                        &node,
674                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
675                    )
676                    .await?;
677                    self.delete_graph_node(&node, parent_scope).await?;
678                    return Ok(node);
679                }
680            }
681
682            let descriptor = self
683                .data_service
684                .metadata
685                .context
686                .require_entity(&node.entity)
687                .map_err(DataServiceError::Runtime)?;
688            let Some(id_property) = descriptor.id_property() else {
689                return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
690                    "entity {} has no id property for graph upsert",
691                    node.entity
692                ))));
693            };
694            let Some(id) = node
695                .values
696                .get(&id_property.name)
697                .filter(|value| !is_unassigned_id_value(value))
698                .cloned()
699            else {
700                // Strip comment to prevent duplicate scope — already captured in active_scope
701                node.comment = None;
702                return self.insert_graph_node_scoped(node, active_scope).await;
703            };
704
705            if node.operation == GraphOperation::Create
706                || self
707                    .fetch_graph_current_row_internal(
708                        &node.entity,
709                        &id_property.name,
710                        &id,
711                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
712                    )
713                    .await?
714                    .is_none()
715            {
716                node.comment = None;
717                return self.insert_graph_node_scoped(node, active_scope).await;
718            }
719
720            let mut one_relations = Vec::new();
721            let mut many_relations = Vec::new();
722            for (name, children) in std::mem::take(&mut node.relations) {
723                let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
724                    DataServiceError::Runtime(RuntimeError::MissingRelation {
725                        entity: node.entity.clone(),
726                        relation: name.clone(),
727                    })
728                })?;
729                match relation.many {
730                    true => many_relations.push((name, relation.clone(), children)),
731                    false => one_relations.push((name, relation.clone(), children)),
732                }
733            }
734
735            for (name, relation, children) in one_relations {
736                if children.len() > 1 {
737                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
738                        "relation {}.{} expects one child, got {}",
739                        node.entity,
740                        name,
741                        children.len()
742                    ))));
743                }
744                let mut saved_children = Vec::new();
745                for child in children {
746                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
747                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
748                    let saved_child = child_repo
749                        .upsert_graph_node_scoped(child, active_scope)
750                        .await?;
751                    if relation.attach {
752                        let foreign_value = saved_child
753                            .values
754                            .get(&relation.foreign_key)
755                            .cloned()
756                            .ok_or_else(|| {
757                                DataServiceError::Runtime(RuntimeError::Graph(format!(
758                                    "saved child {} missing foreign key {} for relation {}.{}",
759                                    relation.target_entity, relation.foreign_key, node.entity, name
760                                )))
761                            })?;
762                        node.values
763                            .insert(relation.local_key.clone(), foreign_value);
764                    }
765                    saved_children.push(saved_child);
766                }
767                node.relations.insert(name, saved_children);
768            }
769
770            let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
771            if !update.values.is_empty() {
772                let prepared_update = self
773                    .prepare_update_command(&update)
774                    .map_err(DataServiceError::Runtime)?;
775                let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
776                self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
777                    .await?;
778                for (field, value) in &prepared_update.values {
779                    node.values.insert(field.clone(), value.clone());
780                }
781                if let Some(version_property) = descriptor.version_property() {
782                    if let Some(expected_version) = prepared_update.expected_version {
783                        node.values.insert(
784                            version_property.name.clone(),
785                            Value::I64(expected_version + 1),
786                        );
787                    }
788                }
789            }
790
791            for (name, relation, children) in many_relations {
792                let local_value =
793                    node.values
794                        .get(&relation.local_key)
795                        .cloned()
796                        .ok_or_else(|| {
797                            DataServiceError::Runtime(RuntimeError::Graph(format!(
798                                "parent {} missing local key {} for relation {}",
799                                node.entity, relation.local_key, name
800                            )))
801                        })?;
802                let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
803                let child_descriptor = self
804                    .data_service
805                    .metadata
806                    .context
807                    .require_entity(&relation.target_entity)
808                    .map_err(DataServiceError::Runtime)?;
809                let child_id_property = child_descriptor.id_property().ok_or_else(|| {
810                    DataServiceError::Runtime(RuntimeError::Graph(format!(
811                        "entity {} has no id property",
812                        relation.target_entity
813                    )))
814                })?;
815
816                let mut seen = std::collections::BTreeSet::new();
817                let mut saved_children = Vec::new();
818                for mut child in children {
819                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
820                    if relation.attach && child.operation != GraphOperation::Reference {
821                        child
822                            .values
823                            .insert(relation.foreign_key.clone(), local_value.clone());
824                    }
825                    if let Some(child_id) = child
826                        .values
827                        .get(&child_id_property.name)
828                        .filter(|value| !is_unassigned_id_value(value))
829                    {
830                        let key = graph_identity_key(child_id);
831                        if !seen.insert(key.clone()) {
832                            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
833                                "duplicate child id {key} in relation {}.{}",
834                                node.entity, name
835                            ))));
836                        }
837                    }
838                    saved_children.push(
839                        child_repo
840                            .upsert_graph_node_scoped(child, active_scope)
841                            .await?,
842                    );
843                }
844
845                node.relations.insert(name, saved_children);
846            }
847
848            Ok(node)
849        })
850    }
851
852    async fn validate_reference_node(
853        &self,
854        node: GraphNode,
855        trace_chain: Vec<teaql_core::TraceNode>,
856    ) -> Result<GraphNode, DataServiceError<E::Error>> {
857        if !node.relations.is_empty() {
858            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
859                "reference node {} cannot contain child relations",
860                node.entity
861            ))));
862        }
863        let descriptor = self
864            .data_service
865            .metadata
866            .context
867            .require_entity(&node.entity)
868            .map_err(DataServiceError::Runtime)?;
869        let id_property = descriptor.id_property().ok_or_else(|| {
870            DataServiceError::Runtime(RuntimeError::Graph(format!(
871                "entity {} has no id property for graph reference",
872                node.entity
873            )))
874        })?;
875        let id = node
876            .values
877            .get(&id_property.name)
878            .filter(|value| !is_unassigned_id_value(value))
879            .cloned()
880            .ok_or_else(|| {
881                DataServiceError::Runtime(RuntimeError::Graph(format!(
882                    "reference node {} missing id property {}",
883                    node.entity, id_property.name
884                )))
885            })?;
886
887        for field in node.values.keys() {
888            if field == &id_property.name {
889                continue;
890            }
891            if descriptor
892                .version_property()
893                .map(|property| field == &property.name)
894                .unwrap_or(false)
895            {
896                continue;
897            }
898            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
899                "reference node {} cannot carry mutable field {}",
900                node.entity, field
901            ))));
902        }
903
904        let current = self
905            .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
906            .await?
907            .ok_or_else(|| {
908                DataServiceError::Runtime(RuntimeError::Graph(format!(
909                    "reference node {}({}) does not exist",
910                    node.entity,
911                    graph_identity_key(&id)
912                )))
913            })?;
914
915        if let Some(version_property) = descriptor.version_property() {
916            if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
917                if *existing_version < 0 {
918                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
919                        "reference node {}({}) is deleted",
920                        node.entity,
921                        graph_identity_key(&id)
922                    ))));
923                }
924                if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
925                {
926                    if expected_version != existing_version {
927                        println!(
928                            "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
929                            node.entity, expected_version, existing_version
930                        );
931                        return Err(DataServiceError::Runtime(
932                            RuntimeError::OptimisticLockConflict {
933                                entity: node.entity,
934                                id: graph_identity_key(&id),
935                            },
936                        ));
937                    }
938                }
939            }
940        }
941
942        Ok(GraphNode {
943            entity: node.entity,
944            values: current,
945            relations: BTreeMap::new(),
946            operation: GraphOperation::Reference,
947            comment: None,
948            dirty_fields: None,
949            original_values: None,
950        })
951    }
952
953    async fn validate_remove_node(
954        &self,
955        node: &GraphNode,
956        trace_chain: Vec<teaql_core::TraceNode>,
957    ) -> Result<(), DataServiceError<E::Error>> {
958        if !node.relations.is_empty() {
959            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
960                "remove node {} cannot contain child relations",
961                node.entity
962            ))));
963        }
964        let descriptor = self
965            .data_service
966            .metadata
967            .context
968            .require_entity(&node.entity)
969            .map_err(DataServiceError::Runtime)?;
970        let id_property = descriptor.id_property().ok_or_else(|| {
971            DataServiceError::Runtime(RuntimeError::Graph(format!(
972                "entity {} has no id property for graph remove",
973                node.entity
974            )))
975        })?;
976        let id = node
977            .values
978            .get(&id_property.name)
979            .filter(|value| !is_unassigned_id_value(value))
980            .cloned()
981            .ok_or_else(|| {
982                DataServiceError::Runtime(RuntimeError::Graph(format!(
983                    "remove node {} missing id property {}",
984                    node.entity, id_property.name
985                )))
986            })?;
987        let current = self
988            .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
989            .await?
990            .ok_or_else(|| {
991                DataServiceError::Runtime(RuntimeError::Graph(format!(
992                    "remove node {}({}) does not exist",
993                    node.entity,
994                    graph_identity_key(&id)
995                )))
996            })?;
997        if let Some(version_property) = descriptor.version_property() {
998            if let Some(Value::I64(existing_version)) = current.get(&version_property.name) {
999                if *existing_version < 0 {
1000                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1001                        "remove node {}({}) is already deleted",
1002                        node.entity,
1003                        graph_identity_key(&id)
1004                    ))));
1005                }
1006            }
1007        }
1008        Ok(())
1009    }
1010
1011    fn graph_node_from_record(
1012        &self,
1013        entity: &str,
1014        record: Record,
1015    ) -> Result<GraphNode, RuntimeError> {
1016        let descriptor = self.data_service.metadata.context.require_entity(entity)?;
1017        let mut node = GraphNode::new(entity);
1018
1019        for (field, value) in record {
1020            if field == "_comment" {
1021                if let Value::Text(comment) = value {
1022                    node.set_comment(comment);
1023                }
1024                continue;
1025            }
1026            if field == "_dirty_fields" {
1027                if let Value::List(fields) = value {
1028                    let mut dirty = std::collections::BTreeSet::new();
1029                    for f in fields {
1030                        if let Value::Text(t) = f {
1031                            dirty.insert(t);
1032                        }
1033                    }
1034                    node.dirty_fields = Some(dirty);
1035                }
1036                continue;
1037            }
1038            if field == "_original_values" {
1039                if let Value::Object(orig) = value {
1040                    node.original_values = Some(orig);
1041                }
1042                continue;
1043            }
1044            let Some(relation) = descriptor.relation_by_name(&field) else {
1045                node.values.insert(field, value);
1046                continue;
1047            };
1048
1049            match value {
1050                Value::Null => {
1051                    node.relations.entry(field).or_default();
1052                }
1053                Value::Object(record) => {
1054                    let child = self.graph_node_from_record(&relation.target_entity, record)?;
1055                    node.relations.entry(field).or_default().push(child);
1056                }
1057                Value::List(values) => {
1058                    let children = node.relations.entry(field.clone()).or_default();
1059                    for value in values {
1060                        let Value::Object(record) = value else {
1061                            return Err(RuntimeError::Graph(format!(
1062                                "relation {}.{} expects object children, got {:?}",
1063                                entity, field, value
1064                            )));
1065                        };
1066                        children
1067                            .push(self.graph_node_from_record(&relation.target_entity, record)?);
1068                    }
1069                }
1070                other => {
1071                    return Err(RuntimeError::Graph(format!(
1072                        "relation {}.{} expects object/list/null, got {:?}",
1073                        entity, field, other
1074                    )));
1075                }
1076            }
1077        }
1078
1079        Ok(node)
1080    }
1081
1082    fn graph_update_command(
1083        &self,
1084        node: &mut GraphNode,
1085        descriptor: &EntityDescriptor,
1086        id_property: &PropertyDescriptor,
1087        id: &Value,
1088    ) -> Result<UpdateCommand, DataServiceError<E::Error>> {
1089        crate::mark_record_status(&mut node.values, crate::CheckObjectStatus::Update);
1090        let check_result = self
1091            .data_service
1092            .metadata
1093            .context
1094            .check_and_fix_record(&node.entity, &mut node.values);
1095        crate::clear_record_status(&mut node.values);
1096        check_result.map_err(DataServiceError::Runtime)?;
1097
1098        let mut command = UpdateCommand::new(node.entity.clone(), id.clone());
1099        command.old_values = node.original_values.clone();
1100        if let Some(version_property) = descriptor.version_property() {
1101            if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1102                command = command.expected_version(*version);
1103            }
1104        }
1105        // Filter properties by dirty_fields when available (Java-style minimal UPDATE).
1106        // When dirty_fields is Some, only modified fields are included in the SET clause.
1107        // When dirty_fields is None (no tracking), fall back to all fields in node.values.
1108        for property in descriptor.properties.iter().filter(|property| {
1109            !property.is_id
1110                && !property.is_version
1111                && property.name != id_property.name
1112                && match &node.dirty_fields {
1113                    Some(dirty) => dirty.contains(&property.name),
1114                    None => node.values.contains_key(&property.name),
1115                }
1116        }) {
1117            if let Some(value) = node.values.get(&property.name) {
1118                command.values.insert(property.name.clone(), value.clone());
1119            }
1120        }
1121        Ok(command)
1122    }
1123
1124    fn delete_graph_node<'b, 's: 'b>(
1125        &'b self,
1126        node: &'b GraphNode,
1127        parent_scope: Option<&'s ScopedCommentNode<'s>>,
1128    ) -> std::pin::Pin<
1129        Box<dyn std::future::Future<Output = Result<u64, DataServiceError<E::Error>>> + Send + '_>,
1130    > {
1131        Box::pin(async move {
1132            let descriptor = self
1133                .data_service
1134                .metadata
1135                .context
1136                .require_entity(&node.entity)
1137                .map_err(DataServiceError::Runtime)?;
1138            let id_property = descriptor.id_property().ok_or_else(|| {
1139                DataServiceError::Runtime(RuntimeError::Graph(format!(
1140                    "entity {} has no id property for graph remove",
1141                    node.entity
1142                )))
1143            })?;
1144            let id = node
1145                .values
1146                .get(&id_property.name)
1147                .filter(|value| !is_unassigned_id_value(value))
1148                .cloned()
1149                .ok_or_else(|| {
1150                    DataServiceError::Runtime(RuntimeError::Graph(format!(
1151                        "remove node {} missing id property {}",
1152                        node.entity, id_property.name
1153                    )))
1154                })?;
1155            let mut delete = DeleteCommand::new(node.entity.clone(), id);
1156            if let Some(version_property) = descriptor.version_property() {
1157                if let Some(Value::I64(version)) = node.values.get(&version_property.name) {
1158                    delete = delete.expected_version(*version);
1159                }
1160            }
1161
1162            // Create scope node for deletion if parent/node comment is present
1163            let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
1164                parent: parent_scope,
1165                track: teaql_core::TraceNode {
1166                    entity_type: node.entity.clone(),
1167                    entity_id: node.id().and_then(|v| match v {
1168                        Value::U64(n) => Some(*n),
1169                        Value::I64(n) => Some(*n as u64),
1170                        _ => None,
1171                    }),
1172                    comment: c.clone(),
1173                },
1174            });
1175            let active_scope = current_scope.as_ref().or(parent_scope);
1176            let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
1177
1178            self.delete_scoped_internal(&delete, lineage).await
1179        })
1180    }
1181
1182    async fn fetch_graph_children(
1183        &self,
1184        entity: &str,
1185        foreign_key: &str,
1186        parent_value: &Value,
1187        trace_chain: Vec<teaql_core::TraceNode>,
1188    ) -> Result<Vec<Record>, DataServiceError<E::Error>> {
1189        let mut query =
1190            SelectQuery::new(entity).filter(Expr::eq(foreign_key, parent_value.clone()));
1191        query.trace_chain = trace_chain;
1192        self.scoped_data_service_internal(entity.to_owned())
1193            .fetch_all_internal(&query)
1194            .await
1195    }
1196    pub(crate) async fn fetch_graph_current_row_internal(
1197        &self,
1198        entity: &str,
1199        id_property: &str,
1200        id: &teaql_core::Value,
1201        trace_chain: Vec<teaql_core::TraceNode>,
1202    ) -> Result<Option<Record>, DataServiceError<E::Error>> {
1203        let mut query = teaql_core::SelectQuery::new(entity)
1204            .filter(teaql_core::Expr::eq(id_property, id.clone()));
1205        query.trace_chain = trace_chain;
1206        let mut rows = self
1207            .scoped_data_service_internal(entity.to_owned())
1208            .fetch_all_internal(&query)
1209            .await?;
1210        Ok(rows.pop())
1211    }
1212
1213    pub(crate) async fn execute_ledger_plan_internal(
1214        &self,
1215        root: crate::EntityRoot,
1216    ) -> Result<(), DataServiceError<E::Error>> {
1217        let comment = root.get_comment();
1218        let trace_chain = comment
1219            .map(|c| {
1220                vec![teaql_core::TraceNode {
1221                    entity_type: self.entity.clone(),
1222                    entity_id: None,
1223                    comment: c,
1224                }]
1225            })
1226            .unwrap_or_default();
1227
1228        let deleted_keys = root.deleted_keys();
1229        let new_keys = root.new_keys();
1230        let change_set = root.current_change_set();
1231
1232        // 1. Execute Deletes
1233        for key in deleted_keys.iter() {
1234            let id = key.id.clone();
1235            let mut cmd = teaql_core::DeleteCommand::new(&key.entity, id);
1236            if let Some(version) = root.get_original_version(key) {
1237                cmd = cmd.expected_version(version);
1238            }
1239            cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1240            self.delete_internal(&cmd).await?;
1241        }
1242
1243        // 2. Execute Updates and Inserts
1244        let mut update_batches: std::collections::BTreeMap<
1245            (String, String),
1246            Vec<crate::EntityKey>,
1247        > = std::collections::BTreeMap::new();
1248        let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1249            std::collections::BTreeMap::new();
1250
1251        for (key, record) in change_set.changes() {
1252            if deleted_keys.contains(key) {
1253                continue;
1254            }
1255            let mut is_new = new_keys.contains(key);
1256
1257            if !is_new {
1258                let descriptor = self
1259                    .data_service
1260                    .metadata
1261                    .context
1262                    .require_entity(&key.entity)
1263                    .map_err(DataServiceError::Runtime)?;
1264                let id_property = descriptor.id_property().ok_or_else(|| {
1265                    DataServiceError::Runtime(RuntimeError::Graph(format!(
1266                        "entity {} has no id property",
1267                        key.entity
1268                    )))
1269                })?;
1270                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1271                let current_row = self
1272                    .fetch_graph_current_row_internal(
1273                        &key.entity,
1274                        &id_property.name,
1275                        &key.id,
1276                        my_trace,
1277                    )
1278                    .await?;
1279                if current_row.is_none() {
1280                    is_new = true;
1281                }
1282            }
1283
1284            match is_new {
1285                true => {
1286                    insert_batches
1287                        .entry(key.entity.clone())
1288                        .or_default()
1289                        .push(key.clone());
1290                }
1291                false => {
1292                    let mut fields: Vec<String> = record.keys().cloned().collect();
1293                    fields.sort();
1294                    let signature = fields.join(",");
1295                    update_batches
1296                        .entry((key.entity.clone(), signature))
1297                        .or_default()
1298                        .push(key.clone());
1299                }
1300            }
1301        }
1302
1303        let mut insert_order: Vec<String> = insert_batches.keys().cloned().collect();
1304        insert_order.sort();
1305        println!("execute_ledger_plan: insert_batches={:?}", insert_order);
1306
1307        for entity in insert_order {
1308            let keys = insert_batches.get(&entity).unwrap();
1309            let descriptor = self
1310                .data_service
1311                .metadata
1312                .context
1313                .require_entity(&entity)
1314                .map_err(DataServiceError::Runtime)?;
1315            let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.table_name);
1316            let mut traces = Vec::new();
1317            for key in keys {
1318                let record = change_set.changes().get(key).unwrap();
1319                let mut db_record = Record::new();
1320                db_record.insert("id".to_owned(), key.id.clone());
1321                for (field, value) in record {
1322                    db_record.insert(field.clone(), value.clone());
1323                }
1324                crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1325                cmd.batch_values.push(db_record);
1326                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1327                traces.push(my_trace);
1328            }
1329            cmd.trace_chains = traces;
1330            self.execute_prepared_batch_insert(cmd).await?;
1331        }
1332
1333        let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1334        update_order.sort();
1335        println!("execute_ledger_plan: update_batches={:?}", update_order);
1336
1337        for signature in update_order {
1338            let keys = update_batches.get(&signature).unwrap();
1339            let descriptor = self
1340                .data_service
1341                .metadata
1342                .context
1343                .require_entity(&signature.0)
1344                .map_err(DataServiceError::Runtime)?;
1345            let mut update_fields: Vec<String> =
1346                signature.0.split(',').map(|s| s.to_string()).collect();
1347            let mut cmd =
1348                teaql_core::BatchUpdateCommand::new(&descriptor.table_name, update_fields);
1349            let mut traces = Vec::new();
1350            for key in keys {
1351                let record = change_set.changes().get(key).unwrap();
1352                let mut db_record = Record::new();
1353                db_record.insert("id".to_owned(), key.id.clone());
1354                for (field, value) in record {
1355                    db_record.insert(field.clone(), value.clone());
1356                }
1357                crate::data_service::helpers::increment_version(
1358                    &mut db_record,
1359                    descriptor,
1360                    root.get_original_version(key),
1361                );
1362                cmd.batch_values.push(db_record);
1363                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1364                traces.push(my_trace);
1365            }
1366            cmd.trace_chains = traces;
1367            self.execute_prepared_batch_update(cmd).await?;
1368        }
1369
1370        Ok(())
1371    }
1372}