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