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