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