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