Skip to main content

teaql_runtime/data_service/
graph.rs

1// Recursive graph futures expose stable internal boundaries used during generator migration.
2#![allow(dead_code, clippy::type_complexity)]
3
4use std::collections::BTreeMap;
5use std::sync::Arc;
6
7use teaql_core::{
8    DeleteCommand, Entity, EntityDescriptor, InsertCommand, MutationValues, PropertyDescriptor,
9    UpdateCommand, Value,
10};
11
12use crate::entity_status::EntityStatus;
13use crate::{
14    DataServiceError, GraphMutationKind, GraphMutationPlan, GraphNode, GraphOperation,
15    RuntimeError, ScopedCommentNode, TraceScopeToken, sorted_update_fields,
16};
17
18use super::{EntityDataService, helpers::*};
19
20fn recover_trace_or_default(token: &Option<Arc<TraceScopeToken>>) -> Vec<teaql_core::TraceNode> {
21    token
22        .as_ref()
23        .map(|t| t.recover_trace_chain())
24        .unwrap_or_default()
25}
26
27fn resolve_trace_chain(
28    specific: Vec<teaql_core::TraceNode>,
29    fallback: &[teaql_core::TraceNode],
30) -> Vec<teaql_core::TraceNode> {
31    match specific.is_empty() {
32        true => fallback.to_vec(),
33        false => specific,
34    }
35}
36
37impl<'a, E> EntityDataService<'a, E>
38where
39    E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
40{
41    pub(crate) async fn save_graph_internal(
42        &self,
43        node: GraphNode,
44    ) -> Result<GraphNode, DataServiceError<E::Error>> {
45        if node.entity != self.entity {
46            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
47                "entity data service {} cannot save graph root {}",
48                self.entity, node.entity
49            ))));
50        }
51        let plan = self.plan_graph(node).await?;
52        self.execute_graph_plan_internal(plan).await
53    }
54
55    pub(crate) async fn save_entity_graph_from_internal(
56        &self,
57        graph: teaql_core::EntityGraph,
58    ) -> Result<GraphNode, DataServiceError<E::Error>> {
59        fn convert(node: teaql_core::EntityGraphNode) -> GraphNode {
60            let mut relations = BTreeMap::new();
61            for (rel_name, child) in node.children {
62                relations
63                    .entry(rel_name)
64                    .or_insert_with(Vec::new)
65                    .push(convert(child));
66            }
67            GraphNode {
68                entity: node.entity_type,
69                values: node.values.into(),
70                relations,
71                operation: match node.operation {
72                    teaql_core::EntityGraphOperation::Save => crate::GraphOperation::Upsert,
73                    teaql_core::EntityGraphOperation::Delete => crate::GraphOperation::Remove,
74                },
75                comment: node.comment,
76                dirty_fields: None,
77                original_values: None,
78            }
79        }
80        self.save_graph_internal(convert(graph.root)).await
81    }
82
83    pub(crate) async fn save_entity_graph_internal<T>(
84        &self,
85        entity: T,
86    ) -> Result<GraphNode, DataServiceError<E::Error>>
87    where
88        T: Entity,
89    {
90        let node = self
91            .graph_node_from_entity(entity)
92            .map_err(DataServiceError::Runtime)?;
93        self.save_graph_internal(node).await
94    }
95
96    pub(crate) async fn save_entity_internal<T>(
97        &self,
98        entity: T,
99        status: EntityStatus,
100    ) -> Result<GraphNode, DataServiceError<E::Error>>
101    where
102        T: Entity,
103    {
104        if !status.need_persist() {
105            return Ok(GraphNode::new(&self.entity));
106        }
107        if status.is_deleted() {
108            let mut node = self
109                .graph_node_from_entity(entity)
110                .map_err(DataServiceError::Runtime)?;
111            node.operation = GraphOperation::Remove;
112            node.relations.clear();
113            return self.save_graph_internal(node).await;
114        }
115        self.save_entity_graph_internal(entity).await
116    }
117    pub(crate) async fn save_entity_with_comment_internal<T>(
118        &self,
119        entity: T,
120        status: EntityStatus,
121        comment: impl Into<String>,
122    ) -> Result<GraphNode, DataServiceError<E::Error>>
123    where
124        T: Entity,
125    {
126        if status.is_deleted() {
127            let mut node = self
128                .graph_node_from_entity(entity)
129                .map_err(DataServiceError::Runtime)?;
130            node.operation = GraphOperation::Remove;
131            node.relations.clear();
132            node.set_comment(comment);
133            return self.save_graph_internal(node).await;
134        }
135        self.save_entity_graph_with_comment_internal(entity, comment)
136            .await
137    }
138    pub(crate) async fn save_entity_graph_with_comment_internal<T>(
139        &self,
140        entity: T,
141        comment: impl Into<String>,
142    ) -> Result<GraphNode, DataServiceError<E::Error>>
143    where
144        T: Entity,
145    {
146        let mut node = self
147            .graph_node_from_entity(entity)
148            .map_err(DataServiceError::Runtime)?;
149        node.set_comment(comment);
150        self.save_graph_internal(node).await
151    }
152
153    /// Create a new entity graph with an annotation comment on the root node.
154    /// This assumes all new nodes do not exist in the database, skipping existence checks
155    /// and throwing an exception on primary key conflict.
156    pub(crate) async fn create_entity_graph_with_comment_internal<T>(
157        &self,
158        entity: T,
159        comment: impl Into<String>,
160    ) -> Result<GraphNode, DataServiceError<E::Error>>
161    where
162        T: Entity,
163    {
164        let mut node = self
165            .graph_node_from_entity(entity)
166            .map_err(DataServiceError::Runtime)?;
167        node.operation = GraphOperation::Create;
168        node.set_comment(comment);
169        self.save_graph_internal(node).await
170    }
171
172    pub async fn plan_graph(
173        &self,
174        node: GraphNode,
175    ) -> Result<GraphMutationPlan, DataServiceError<E::Error>> {
176        if node.entity != self.entity {
177            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
178                "entity data service {} cannot plan graph root {}",
179                self.entity, node.entity
180            ))));
181        }
182        let mut node = node;
183        let mut plan = GraphMutationPlan::default();
184        self.collect_graph_plan(&mut node, &mut plan, None, None, false)
185            .await?;
186        plan.planned_root = Some(node);
187        plan.rebuild_batches();
188        Ok(plan)
189    }
190
191    pub(crate) async fn execute_graph_plan_internal(
192        &self,
193        plan: GraphMutationPlan,
194    ) -> Result<GraphNode, DataServiceError<E::Error>> {
195        let Some(mut root) = plan.planned_root else {
196            return Err(DataServiceError::Runtime(RuntimeError::Graph(
197                "graph mutation plan has no planned root".to_owned(),
198            )));
199        };
200
201        for batch in plan.batches {
202            if batch.items.is_empty()
203                || (matches!(batch.kind, GraphMutationKind::Update)
204                    && batch.update_fields.is_empty())
205            {
206                continue;
207            }
208            match batch.kind {
209                GraphMutationKind::Create => {
210                    let mut cmd = teaql_core::BatchInsertCommand::new(&batch.entity);
211                    for item in batch.items {
212                        cmd.batch_values.push(item.values);
213                        cmd.trace_chains
214                            .push(recover_trace_or_default(&item.scope_token));
215                    }
216                    self.execute_prepared_batch_insert(cmd).await?;
217                }
218                GraphMutationKind::Update => {
219                    if batch.update_fields.is_empty() {
220                        continue;
221                    }
222                    let mut cmd =
223                        teaql_core::BatchUpdateCommand::new(&batch.entity, batch.update_fields);
224                    for item in batch.items {
225                        let id = item.values.get("id").cloned().ok_or_else(|| {
226                            DataServiceError::Runtime(RuntimeError::Graph(format!(
227                                "update item in batch missing id for {}",
228                                batch.entity
229                            )))
230                        })?;
231                        let version = item.values.get("version").and_then(|v| match v {
232                            teaql_core::Value::I64(n) => Some(*n),
233                            _ => None,
234                        });
235                        cmd.batch_values.push(item.values);
236                        cmd.batch_ids.push(id);
237                        cmd.batch_expected_versions.push(version);
238                        cmd.batch_old_values.push(item.old_values);
239                        cmd.trace_chains
240                            .push(recover_trace_or_default(&item.scope_token));
241                    }
242                    self.execute_prepared_batch_update(cmd).await?;
243                }
244                GraphMutationKind::Delete => {
245                    // For now, loop individually since we lack BatchDeleteCommand
246                    for item in batch.items {
247                        let id = item.values.get("id").cloned().ok_or_else(|| {
248                            DataServiceError::Runtime(RuntimeError::Graph(format!(
249                                "delete item in batch missing id for {}",
250                                batch.entity
251                            )))
252                        })?;
253                        let mut cmd = teaql_core::DeleteCommand::new(&batch.entity, id);
254                        if let Some(teaql_core::Value::I64(version)) = item.values.get("version") {
255                            cmd = cmd.expected_version(*version);
256                        }
257                        let trace_chain = recover_trace_or_default(&item.scope_token);
258                        self.delete_scoped_internal(&cmd, trace_chain).await?;
259                    }
260                }
261                GraphMutationKind::Reference => {
262                    // References are skipped in execution, they only validate during traversal
263                }
264            }
265        }
266
267        if root.operation != GraphOperation::Remove {
268            let descriptor = self
269                .data_service
270                .metadata
271                .context
272                .require_entity(&root.entity)
273                .map_err(DataServiceError::Runtime)?;
274            let id_property = descriptor.id_property().ok_or_else(|| {
275                DataServiceError::Runtime(RuntimeError::Graph(format!(
276                    "entity {} has no id property",
277                    root.entity
278                )))
279            })?;
280            let id = root.values.get(&id_property.name).cloned().ok_or_else(|| {
281                DataServiceError::Runtime(RuntimeError::Graph(format!(
282                    "saved {} missing identity field {}",
283                    root.entity, id_property.name
284                )))
285            })?;
286            root.values = self
287                .fetch_graph_current_row_internal(
288                    &root.entity,
289                    &id_property.name,
290                    &id,
291                    root.comment
292                        .clone()
293                        .map(|comment| {
294                            vec![teaql_core::TraceNode {
295                                kind: teaql_core::TraceKind::AuditReason,
296                                entity_type: root.entity.clone(),
297                                entity_id: id.try_u64(),
298                                comment,
299                            }]
300                        })
301                        .unwrap_or_default(),
302                )
303                .await?
304                .map(Into::into)
305                .ok_or_else(|| {
306                    DataServiceError::Runtime(RuntimeError::Graph(format!(
307                        "persisted {} record could not be read back",
308                        root.entity
309                    )))
310                })?;
311        }
312
313        Ok(root)
314    }
315
316    pub fn graph_node_from_entity<T>(&self, entity: T) -> Result<GraphNode, RuntimeError>
317    where
318        T: Entity,
319    {
320        let descriptor = T::entity_descriptor();
321        if descriptor.name != self.entity {
322            return Err(RuntimeError::Graph(format!(
323                "entity data service {} cannot extract graph root {}",
324                self.entity, descriptor.name
325            )));
326        }
327        // Extract dirty field names before into_values() consumes the entity.
328        // This is the Rust equivalent of Java's entity.getUpdatedProperties().
329        let dirty_fields = entity.dirty_fields();
330        let original_values = entity.original_values();
331        let is_deleted = entity.is_marked_as_delete();
332        let comment = entity.get_comment();
333        let mut node = self.graph_node_from_values(&descriptor.name, entity.into_values())?;
334        node.dirty_fields = dirty_fields;
335        node.original_values = original_values;
336        if is_deleted {
337            node.operation = GraphOperation::Remove;
338            node.relations.clear();
339        }
340        if let Some(c) = comment {
341            node.set_comment(c);
342        }
343        Ok(node)
344    }
345
346    fn collect_graph_plan<'b, 's: 'b>(
347        &'b self,
348        node: &'b mut GraphNode,
349        plan: &'b mut GraphMutationPlan,
350        parent_scope: Option<&'s ScopedCommentNode<'s>>,
351        parent_token: Option<Arc<TraceScopeToken>>,
352        parent_is_create: bool,
353    ) -> std::pin::Pin<
354        Box<dyn std::future::Future<Output = Result<(), DataServiceError<E::Error>>> + Send + '_>,
355    > {
356        Box::pin(async move {
357            match node.operation {
358                GraphOperation::Reference => {
359                    plan.push(
360                        node.entity.clone(),
361                        GraphMutationKind::Reference,
362                        node.values.clone().into(),
363                        Vec::new(),
364                        parent_token,
365                        node.original_values.clone(),
366                    );
367                    return Ok(());
368                }
369                GraphOperation::Remove => {
370                    plan.push(
371                        node.entity.clone(),
372                        GraphMutationKind::Delete,
373                        node.values.clone().into(),
374                        Vec::new(),
375                        parent_token,
376                        node.original_values.clone(),
377                    );
378                    return Ok(());
379                }
380                GraphOperation::Upsert | GraphOperation::Create => {}
381            }
382
383            let descriptor = self
384                .data_service
385                .metadata
386                .context
387                .require_entity(&node.entity)
388                .map_err(DataServiceError::Runtime)?;
389
390            // Create scope node on the current stack frame if this node has a comment
391            let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
392                parent: parent_scope,
393                track: teaql_core::TraceNode {
394                    kind: teaql_core::TraceKind::AuditReason,
395                    entity_type: node.entity.clone(),
396                    entity_id: node.id().and_then(|v| match v {
397                        Value::U64(n) => Some(*n),
398                        Value::I64(n) => Some(*n as u64),
399                        _ => None,
400                    }),
401                    comment: c.clone(),
402                },
403            });
404            let active_scope = current_scope.as_ref().or(parent_scope);
405
406            let id_property = descriptor.id_property().cloned();
407            let id = id_property.as_ref().and_then(|property| {
408                node.values
409                    .get(&property.name)
410                    .filter(|value| !is_unassigned_id_value(value))
411                    .cloned()
412            });
413
414            if let Some(id_val) = &id
415                && !plan
416                    .visited_nodes
417                    .insert((node.entity.clone(), graph_identity_key(id_val)))
418            {
419                return Ok(());
420            }
421
422            let is_create_op = node.operation == GraphOperation::Create
423                || (parent_is_create && node.operation == GraphOperation::Upsert);
424
425            let is_update = match is_create_op {
426                true => false,
427                false => match (id_property.as_ref(), id.as_ref()) {
428                    (Some(id_property), Some(id)) => self
429                        .fetch_graph_current_row_internal(
430                            &node.entity,
431                            &id_property.name,
432                            id,
433                            active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
434                        )
435                        .await?
436                        .is_some(),
437                    _ => false,
438                },
439            };
440            if !is_update {
441                if let Some(id_property) = id_property.as_ref() {
442                    let needs_id = !node.values.contains_key(&id_property.name)
443                        || node
444                            .values
445                            .get(&id_property.name)
446                            .is_some_and(is_unassigned_id_value);
447                    if needs_id {
448                        let id = self
449                            .data_service
450                            .metadata
451                            .context
452                            .next_id(&node.entity)
453                            .map_err(DataServiceError::Runtime)?;
454                        node.values.insert(id_property.name.clone(), Value::U64(id));
455                    }
456                }
457                ensure_initial_version(&mut node.values, descriptor);
458                crate::data_service::helpers::ensure_timestamps(&mut node.values, descriptor, true);
459            } else {
460                crate::data_service::helpers::ensure_timestamps(
461                    &mut node.values,
462                    descriptor,
463                    false,
464                );
465            }
466            let update_fields = if is_update {
467                let mut excluded = Vec::new();
468                if let Some(id_property) = id_property.as_ref() {
469                    excluded.push(id_property.name.clone());
470                }
471                if let Some(version_property) = descriptor.version_property() {
472                    excluded.push(version_property.name.clone());
473                }
474                let mut fields = sorted_update_fields(&node.values, excluded);
475                if let Some(dirty) = &node.dirty_fields {
476                    fields.retain(|f| dirty.contains(f));
477                }
478                fields
479            } else {
480                Default::default()
481            };
482
483            // Build the TraceScopeToken for this node (only if it has a comment).
484            // This is an Arc-linked persistent list: zero-copy, O(1) creation.
485            let current_token = node
486                .comment
487                .as_ref()
488                .map(|c| {
489                    Arc::new(TraceScopeToken {
490                        parent: parent_token.clone(),
491                        track: teaql_core::TraceNode {
492                            kind: teaql_core::TraceKind::AuditReason,
493                            entity_type: node.entity.clone(),
494                            entity_id: node.id().and_then(|v| match v {
495                                Value::U64(n) => Some(*n),
496                                Value::I64(n) => Some(*n as u64),
497                                _ => None,
498                            }),
499                            comment: c.clone(),
500                        },
501                        node_index: plan.next_item_index,
502                    })
503                })
504                .or_else(|| parent_token.clone());
505
506            plan.push(
507                node.entity.clone(),
508                GraphMutationKind::for_update(is_update),
509                node.values.clone().into(),
510                update_fields,
511                current_token.clone(),
512                node.original_values.clone(),
513            );
514
515            for (name, children) in &mut node.relations {
516                let relation = descriptor.relation_by_name(name).ok_or_else(|| {
517                    DataServiceError::Runtime(RuntimeError::MissingRelation {
518                        entity: node.entity.clone(),
519                        relation: name.clone(),
520                    })
521                })?;
522                let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
523                for child in children {
524                    ensure_relation_target(&node.entity, name, &relation.target_entity, child)?;
525                    child_repo
526                        .collect_graph_plan(
527                            child,
528                            plan,
529                            active_scope,
530                            current_token.clone(),
531                            is_create_op,
532                        )
533                        .await?;
534                }
535            }
536            Ok(())
537        })
538    }
539
540    fn insert_graph_node_scoped<'b, 's: 'b>(
541        &'b self,
542        mut node: GraphNode,
543        parent_scope: Option<&'s ScopedCommentNode<'s>>,
544    ) -> std::pin::Pin<
545        Box<
546            dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
547                + Send
548                + '_,
549        >,
550    > {
551        Box::pin(async move {
552            match node.operation {
553                GraphOperation::Upsert | GraphOperation::Create => {}
554                GraphOperation::Reference => {
555                    return self
556                        .validate_reference_node(
557                            node,
558                            parent_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
559                        )
560                        .await;
561                }
562                GraphOperation::Remove => {
563                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
564                        "create graph cannot remove node {}",
565                        node.entity
566                    ))));
567                }
568            }
569
570            // Create scope node on the current stack frame if this node has a comment
571            let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
572                parent: parent_scope,
573                track: teaql_core::TraceNode {
574                    kind: teaql_core::TraceKind::AuditReason,
575                    entity_type: node.entity.clone(),
576                    entity_id: node.id().and_then(|v| match v {
577                        Value::U64(n) => Some(*n),
578                        Value::I64(n) => Some(*n as u64),
579                        _ => None,
580                    }),
581                    comment: c.clone(),
582                },
583            });
584            let active_scope = current_scope.as_ref().or(parent_scope);
585
586            let descriptor = self
587                .data_service
588                .metadata
589                .context
590                .require_entity(&node.entity)
591                .map_err(DataServiceError::Runtime)?;
592
593            let mut one_relations = Vec::new();
594            let mut many_relations = Vec::new();
595            for (name, children) in std::mem::take(&mut node.relations) {
596                let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
597                    DataServiceError::Runtime(RuntimeError::MissingRelation {
598                        entity: node.entity.clone(),
599                        relation: name.clone(),
600                    })
601                })?;
602                match relation.many {
603                    true => many_relations.push((name, relation.clone(), children)),
604                    false => one_relations.push((name, relation.clone(), children)),
605                }
606            }
607
608            for (name, relation, children) in one_relations {
609                if children.len() > 1 {
610                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
611                        "relation {}.{} expects one child, got {}",
612                        node.entity,
613                        name,
614                        children.len()
615                    ))));
616                }
617                let mut saved_children = Vec::new();
618                for child in children {
619                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
620                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
621                    let saved_child = child_repo
622                        .insert_graph_node_scoped(child, active_scope)
623                        .await?;
624                    if relation.attach {
625                        let foreign_value = saved_child
626                            .values
627                            .get(&relation.foreign_key)
628                            .cloned()
629                            .ok_or_else(|| {
630                                DataServiceError::Runtime(RuntimeError::Graph(format!(
631                                    "saved child {} missing foreign key {} for relation {}.{}",
632                                    relation.target_entity, relation.foreign_key, node.entity, name
633                                )))
634                            })?;
635                        node.values
636                            .insert(relation.local_key.clone(), foreign_value);
637                    }
638                    saved_children.push(saved_child);
639                }
640                node.relations.insert(name, saved_children);
641            }
642
643            let command = self
644                .prepare_insert_command(&InsertCommand {
645                    entity: node.entity.clone(),
646                    values: node.values.clone().into(),
647                    trace_chain: Vec::new(),
648                })
649                .map_err(DataServiceError::Runtime)?;
650            let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
651            self.execute_prepared_insert_with_comment(command.clone(), lineage)
652                .await?;
653            node.values = command.values.into();
654            if let Some(id_property) = descriptor.id_property()
655                && let Some(id) = node.values.get(&id_property.name).cloned()
656            {
657                node.values = self
658                    .fetch_graph_current_row_internal(
659                        &node.entity,
660                        &id_property.name,
661                        &id,
662                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
663                    )
664                    .await?
665                    .map(Into::into)
666                    .ok_or_else(|| {
667                        DataServiceError::Runtime(RuntimeError::Graph(format!(
668                            "persisted {} record could not be read back",
669                            node.entity
670                        )))
671                    })?;
672            }
673
674            for (name, relation, children) in many_relations {
675                let local_value =
676                    node.values
677                        .get(&relation.local_key)
678                        .cloned()
679                        .ok_or_else(|| {
680                            DataServiceError::Runtime(RuntimeError::Graph(format!(
681                                "parent {} missing local key {} for relation {}",
682                                node.entity, relation.local_key, name
683                            )))
684                        })?;
685                let mut saved_children = Vec::new();
686                for mut child in children {
687                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
688                    if relation.attach {
689                        child
690                            .values
691                            .insert(relation.foreign_key.clone(), local_value.clone());
692                    }
693                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
694                    saved_children.push(
695                        child_repo
696                            .insert_graph_node_scoped(child, active_scope)
697                            .await?,
698                    );
699                }
700                node.relations.insert(name, saved_children);
701            }
702
703            Ok(node)
704        })
705    }
706
707    fn upsert_graph_node_scoped<'b, 's: 'b>(
708        &'b self,
709        mut node: GraphNode,
710        parent_scope: Option<&'s ScopedCommentNode<'s>>,
711    ) -> std::pin::Pin<
712        Box<
713            dyn std::future::Future<Output = Result<GraphNode, DataServiceError<E::Error>>>
714                + Send
715                + '_,
716        >,
717    > {
718        Box::pin(async move {
719            // Create scope node on the current stack frame if this node has a comment
720            let current_scope = node.comment.as_ref().map(|c| ScopedCommentNode {
721                parent: parent_scope,
722                track: teaql_core::TraceNode {
723                    kind: teaql_core::TraceKind::AuditReason,
724                    entity_type: node.entity.clone(),
725                    entity_id: node.id().and_then(|v| match v {
726                        Value::U64(n) => Some(*n),
727                        Value::I64(n) => Some(*n as u64),
728                        _ => None,
729                    }),
730                    comment: c.clone(),
731                },
732            });
733            let active_scope = current_scope.as_ref().or(parent_scope);
734
735            match node.operation {
736                GraphOperation::Upsert | GraphOperation::Create => {}
737                GraphOperation::Reference => {
738                    return self
739                        .validate_reference_node(
740                            node,
741                            active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
742                        )
743                        .await;
744                }
745                GraphOperation::Remove => {
746                    self.validate_remove_node(
747                        &node,
748                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
749                    )
750                    .await?;
751                    self.delete_graph_node(&node, parent_scope).await?;
752                    return Ok(node);
753                }
754            }
755
756            let descriptor = self
757                .data_service
758                .metadata
759                .context
760                .require_entity(&node.entity)
761                .map_err(DataServiceError::Runtime)?;
762            let Some(id_property) = descriptor.id_property() else {
763                return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
764                    "entity {} has no id property for graph upsert",
765                    node.entity
766                ))));
767            };
768            let Some(id) = node
769                .values
770                .get(&id_property.name)
771                .filter(|value| !is_unassigned_id_value(value))
772                .cloned()
773            else {
774                // Strip comment to prevent duplicate scope — already captured in active_scope
775                node.comment = None;
776                return self.insert_graph_node_scoped(node, active_scope).await;
777            };
778
779            if node.operation == GraphOperation::Create
780                || self
781                    .fetch_graph_current_row_internal(
782                        &node.entity,
783                        &id_property.name,
784                        &id,
785                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
786                    )
787                    .await?
788                    .is_none()
789            {
790                node.comment = None;
791                return self.insert_graph_node_scoped(node, active_scope).await;
792            }
793
794            let mut one_relations = Vec::new();
795            let mut many_relations = Vec::new();
796            for (name, children) in std::mem::take(&mut node.relations) {
797                let relation = descriptor.relation_by_name(&name).ok_or_else(|| {
798                    DataServiceError::Runtime(RuntimeError::MissingRelation {
799                        entity: node.entity.clone(),
800                        relation: name.clone(),
801                    })
802                })?;
803                match relation.many {
804                    true => many_relations.push((name, relation.clone(), children)),
805                    false => one_relations.push((name, relation.clone(), children)),
806                }
807            }
808
809            for (name, relation, children) in one_relations {
810                if children.len() > 1 {
811                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
812                        "relation {}.{} expects one child, got {}",
813                        node.entity,
814                        name,
815                        children.len()
816                    ))));
817                }
818                let mut saved_children = Vec::new();
819                for child in children {
820                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
821                    let child_repo = self.scoped_data_service_internal(child.entity.clone());
822                    let saved_child = child_repo
823                        .upsert_graph_node_scoped(child, active_scope)
824                        .await?;
825                    if relation.attach {
826                        let foreign_value = saved_child
827                            .values
828                            .get(&relation.foreign_key)
829                            .cloned()
830                            .ok_or_else(|| {
831                                DataServiceError::Runtime(RuntimeError::Graph(format!(
832                                    "saved child {} missing foreign key {} for relation {}.{}",
833                                    relation.target_entity, relation.foreign_key, node.entity, name
834                                )))
835                            })?;
836                        node.values
837                            .insert(relation.local_key.clone(), foreign_value);
838                    }
839                    saved_children.push(saved_child);
840                }
841                node.relations.insert(name, saved_children);
842            }
843
844            let update = self.graph_update_command(&mut node, descriptor, id_property, &id)?;
845            if !update.values.is_empty() {
846                let prepared_update = self
847                    .prepare_update_command(&update)
848                    .map_err(DataServiceError::Runtime)?;
849                let lineage = active_scope.map(|s| s.to_trace_chain()).unwrap_or_default();
850                self.execute_prepared_update_with_comment(prepared_update.clone(), lineage)
851                    .await?;
852                for (field, value) in &prepared_update.values {
853                    node.values.insert(field.clone(), value.clone());
854                }
855                if let Some(version_property) = descriptor.version_property()
856                    && let Some(expected_version) = prepared_update.expected_version
857                {
858                    node.values.insert(
859                        version_property.name.clone(),
860                        Value::I64(expected_version + 1),
861                    );
862                }
863                node.values = self
864                    .fetch_graph_current_row_internal(
865                        &node.entity,
866                        &id_property.name,
867                        &id,
868                        active_scope.map(|s| s.to_trace_chain()).unwrap_or_default(),
869                    )
870                    .await?
871                    .map(Into::into)
872                    .ok_or_else(|| {
873                        DataServiceError::Runtime(RuntimeError::Graph(format!(
874                            "persisted {} record could not be read back",
875                            node.entity
876                        )))
877                    })?;
878            }
879
880            for (name, relation, children) in many_relations {
881                let local_value =
882                    node.values
883                        .get(&relation.local_key)
884                        .cloned()
885                        .ok_or_else(|| {
886                            DataServiceError::Runtime(RuntimeError::Graph(format!(
887                                "parent {} missing local key {} for relation {}",
888                                node.entity, relation.local_key, name
889                            )))
890                        })?;
891                let child_repo = self.scoped_data_service_internal(relation.target_entity.clone());
892                let child_descriptor = self
893                    .data_service
894                    .metadata
895                    .context
896                    .require_entity(&relation.target_entity)
897                    .map_err(DataServiceError::Runtime)?;
898                let child_id_property = child_descriptor.id_property().ok_or_else(|| {
899                    DataServiceError::Runtime(RuntimeError::Graph(format!(
900                        "entity {} has no id property",
901                        relation.target_entity
902                    )))
903                })?;
904
905                let mut seen = std::collections::BTreeSet::new();
906                let mut saved_children = Vec::new();
907                for mut child in children {
908                    ensure_relation_target(&node.entity, &name, &relation.target_entity, &child)?;
909                    if relation.attach && child.operation != GraphOperation::Reference {
910                        child
911                            .values
912                            .insert(relation.foreign_key.clone(), local_value.clone());
913                    }
914                    if let Some(child_id) = child
915                        .values
916                        .get(&child_id_property.name)
917                        .filter(|value| !is_unassigned_id_value(value))
918                    {
919                        let key = graph_identity_key(child_id);
920                        if !seen.insert(key.clone()) {
921                            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
922                                "duplicate child id {key} in relation {}.{}",
923                                node.entity, name
924                            ))));
925                        }
926                    }
927                    saved_children.push(
928                        child_repo
929                            .upsert_graph_node_scoped(child, active_scope)
930                            .await?,
931                    );
932                }
933
934                node.relations.insert(name, saved_children);
935            }
936
937            Ok(node)
938        })
939    }
940
941    async fn validate_reference_node(
942        &self,
943        node: GraphNode,
944        trace_chain: Vec<teaql_core::TraceNode>,
945    ) -> Result<GraphNode, DataServiceError<E::Error>> {
946        if !node.relations.is_empty() {
947            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
948                "reference node {} cannot contain child relations",
949                node.entity
950            ))));
951        }
952        let descriptor = self
953            .data_service
954            .metadata
955            .context
956            .require_entity(&node.entity)
957            .map_err(DataServiceError::Runtime)?;
958        let id_property = descriptor.id_property().ok_or_else(|| {
959            DataServiceError::Runtime(RuntimeError::Graph(format!(
960                "entity {} has no id property for graph reference",
961                node.entity
962            )))
963        })?;
964        let id = node
965            .values
966            .get(&id_property.name)
967            .filter(|value| !is_unassigned_id_value(value))
968            .cloned()
969            .ok_or_else(|| {
970                DataServiceError::Runtime(RuntimeError::Graph(format!(
971                    "reference node {} missing id property {}",
972                    node.entity, id_property.name
973                )))
974            })?;
975
976        for field in node.values.keys() {
977            if field == &id_property.name {
978                continue;
979            }
980            if descriptor
981                .version_property()
982                .map(|property| field == &property.name)
983                .unwrap_or(false)
984            {
985                continue;
986            }
987            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
988                "reference node {} cannot carry mutable field {}",
989                node.entity, field
990            ))));
991        }
992
993        let current = self
994            .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
995            .await?
996            .ok_or_else(|| {
997                DataServiceError::Runtime(RuntimeError::Graph(format!(
998                    "reference node {}({}) does not exist",
999                    node.entity,
1000                    graph_identity_key(&id)
1001                )))
1002            })?;
1003
1004        if let Some(version_property) = descriptor.version_property()
1005            && let Some(Value::I64(existing_version)) = current.get(&version_property.name)
1006        {
1007            if *existing_version < 0 {
1008                return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1009                    "reference node {}({}) is deleted",
1010                    node.entity,
1011                    graph_identity_key(&id)
1012                ))));
1013            }
1014            if let Some(Value::I64(expected_version)) = node.values.get(&version_property.name)
1015                && expected_version != existing_version
1016            {
1017                println!(
1018                    "OptimisticLockConflict in validate_reference_node! entity={}, expected={}, existing={}",
1019                    node.entity, expected_version, existing_version
1020                );
1021                return Err(DataServiceError::Runtime(
1022                    RuntimeError::OptimisticLockConflict {
1023                        entity: node.entity,
1024                        id: graph_identity_key(&id),
1025                    },
1026                ));
1027            }
1028        }
1029
1030        Ok(GraphNode {
1031            entity: node.entity,
1032            values: current.into(),
1033            relations: BTreeMap::new(),
1034            operation: GraphOperation::Reference,
1035            comment: None,
1036            dirty_fields: None,
1037            original_values: None,
1038        })
1039    }
1040
1041    async fn validate_remove_node(
1042        &self,
1043        node: &GraphNode,
1044        trace_chain: Vec<teaql_core::TraceNode>,
1045    ) -> Result<(), DataServiceError<E::Error>> {
1046        if !node.relations.is_empty() {
1047            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1048                "remove node {} cannot contain child relations",
1049                node.entity
1050            ))));
1051        }
1052        let descriptor = self
1053            .data_service
1054            .metadata
1055            .context
1056            .require_entity(&node.entity)
1057            .map_err(DataServiceError::Runtime)?;
1058        let id_property = descriptor.id_property().ok_or_else(|| {
1059            DataServiceError::Runtime(RuntimeError::Graph(format!(
1060                "entity {} has no id property for graph remove",
1061                node.entity
1062            )))
1063        })?;
1064        let id = node
1065            .values
1066            .get(&id_property.name)
1067            .filter(|value| !is_unassigned_id_value(value))
1068            .cloned()
1069            .ok_or_else(|| {
1070                DataServiceError::Runtime(RuntimeError::Graph(format!(
1071                    "remove node {} missing id property {}",
1072                    node.entity, id_property.name
1073                )))
1074            })?;
1075        let current = self
1076            .fetch_graph_current_row_internal(&node.entity, &id_property.name, &id, trace_chain)
1077            .await?
1078            .ok_or_else(|| {
1079                DataServiceError::Runtime(RuntimeError::Graph(format!(
1080                    "remove node {}({}) does not exist",
1081                    node.entity,
1082                    graph_identity_key(&id)
1083                )))
1084            })?;
1085        if let Some(version_property) = descriptor.version_property()
1086            && let Some(Value::I64(existing_version)) = current.get(&version_property.name)
1087            && *existing_version < 0
1088        {
1089            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1090                "remove node {}({}) is already deleted",
1091                node.entity,
1092                graph_identity_key(&id)
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();
1201        if let Some(version_property) = descriptor.version_property()
1202            && let Some(Value::I64(version)) = node.values.get(&version_property.name)
1203        {
1204            command = command.expected_version(*version);
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                && let Some(Value::I64(version)) = node.values.get(&version_property.name)
1259            {
1260                delete = delete.expected_version(*version);
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) fn order_new_ledger_keys(
1302        &self,
1303        keys: impl IntoIterator<Item = crate::EntityKey>,
1304        changes: &std::collections::BTreeMap<crate::EntityKey, crate::EntityValues>,
1305    ) -> Result<Vec<crate::EntityKey>, RuntimeError> {
1306        let keys = keys.into_iter().collect::<std::collections::BTreeSet<_>>();
1307        let mut incoming = keys
1308            .iter()
1309            .cloned()
1310            .map(|key| (key, 0_usize))
1311            .collect::<std::collections::BTreeMap<_, _>>();
1312        let mut outgoing = std::collections::BTreeMap::<
1313            crate::EntityKey,
1314            std::collections::BTreeSet<crate::EntityKey>,
1315        >::new();
1316
1317        let add_dependency =
1318            |parent: &crate::EntityKey,
1319             child: &crate::EntityKey,
1320             incoming: &mut std::collections::BTreeMap<crate::EntityKey, usize>,
1321             outgoing: &mut std::collections::BTreeMap<
1322                crate::EntityKey,
1323                std::collections::BTreeSet<crate::EntityKey>,
1324            >| {
1325                if parent != child
1326                    && outgoing
1327                        .entry(parent.clone())
1328                        .or_default()
1329                        .insert(child.clone())
1330                {
1331                    *incoming
1332                        .get_mut(child)
1333                        .expect("every inserted key has an indegree") += 1;
1334                }
1335            };
1336
1337        for source in &keys {
1338            let descriptor = self
1339                .data_service
1340                .metadata
1341                .context
1342                .require_entity(source.entity.as_ref())?;
1343            let Some(values) = changes.get(source) else {
1344                continue;
1345            };
1346            for relation in &descriptor.relations {
1347                let Some(source_value) = values
1348                    .get(&relation.local_key)
1349                    .filter(|value| !matches!(value, Value::Null | Value::TypedNull(_)))
1350                else {
1351                    continue;
1352                };
1353
1354                for target in keys
1355                    .iter()
1356                    .filter(|target| target.entity.as_ref() == relation.target_entity)
1357                {
1358                    let target_value = changes
1359                        .get(target)
1360                        .and_then(|record| record.get(&relation.foreign_key))
1361                        .or_else(|| (relation.foreign_key == "id").then_some(&target.id));
1362                    if target_value != Some(source_value) {
1363                        continue;
1364                    }
1365                    if relation.many {
1366                        add_dependency(source, target, &mut incoming, &mut outgoing);
1367                    } else {
1368                        add_dependency(target, source, &mut incoming, &mut outgoing);
1369                    }
1370                }
1371            }
1372        }
1373
1374        let mut ready = incoming
1375            .iter()
1376            .filter(|(_, count)| **count == 0)
1377            .map(|(key, _)| key.clone())
1378            .collect::<std::collections::BTreeSet<_>>();
1379        let mut ordered = Vec::with_capacity(keys.len());
1380        while let Some(key) = ready.pop_first() {
1381            ordered.push(key.clone());
1382            if let Some(dependents) = outgoing.get(&key) {
1383                for dependent in dependents {
1384                    let count = incoming
1385                        .get_mut(dependent)
1386                        .expect("every dependent key has an indegree");
1387                    *count -= 1;
1388                    if *count == 0 {
1389                        ready.insert(dependent.clone());
1390                    }
1391                }
1392            }
1393        }
1394
1395        if ordered.len() != keys.len() {
1396            let cyclic_entities = incoming
1397                .into_iter()
1398                .filter(|(_, count)| *count > 0)
1399                .map(|(key, _)| key.entity.into_owned())
1400                .collect::<std::collections::BTreeSet<_>>()
1401                .into_iter()
1402                .collect::<Vec<_>>()
1403                .join(", ");
1404            return Err(RuntimeError::Graph(format!(
1405                "new entity graph contains a required foreign-key cycle among: {cyclic_entities}"
1406            )));
1407        }
1408        Ok(ordered)
1409    }
1410
1411    pub(crate) async fn execute_ledger_plan_internal(
1412        &self,
1413        root: crate::EntityRuntimeState,
1414        locations: &std::collections::BTreeMap<crate::EntityKey, crate::ObjectLocation>,
1415    ) -> Result<std::collections::BTreeMap<crate::EntityKey, Value>, DataServiceError<E::Error>>
1416    {
1417        if let Some(error) = root.first_composition_error() {
1418            return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1419                "generated entity graph attachment failed before save: {error}"
1420            ))));
1421        }
1422        let mut generated_ids = std::collections::BTreeMap::new();
1423        let comment = root.get_comment();
1424        let trace_chain = comment
1425            .map(|c| {
1426                vec![teaql_core::TraceNode {
1427                    kind: teaql_core::TraceKind::AuditReason,
1428                    entity_type: self.entity.clone(),
1429                    entity_id: None,
1430                    comment: c,
1431                }]
1432            })
1433            .unwrap_or_default();
1434
1435        let deleted_keys = root.deleted_keys();
1436        let new_keys = root.new_keys();
1437        let change_set = root.current_change_set();
1438
1439        // Deletion of an existing versioned entity must carry the version
1440        // obtained by loading that entity. A newly created entity deleted
1441        // before save has no database row and is cancelled within the graph.
1442        for key in &deleted_keys {
1443            if new_keys.contains(key) {
1444                continue;
1445            }
1446            let descriptor = self
1447                .data_service
1448                .metadata
1449                .context
1450                .require_entity(&key.entity)
1451                .map_err(DataServiceError::Runtime)?;
1452            if descriptor.version_property().is_some() && root.get_original_version(key).is_none() {
1453                return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1454                    "cannot delete {}({:?}) without its loaded original version; load the full entity before mutation",
1455                    key.entity, key.id
1456                ))));
1457            }
1458        }
1459
1460        // `save_audited_ledger_entity` has already preflighted the complete
1461        // typed graph before entering this executor and merged every Fix value
1462        // back into the ledger.  Do not run Checker/Fix again over these sparse
1463        // change records: unchanged loaded fields are intentionally absent, and
1464        // a second pass both violates once-per-save semantics and misclassifies
1465        // them as NotLoaded.
1466        let mut checked_changes = std::collections::BTreeMap::new();
1467        for (key, record) in change_set.changes() {
1468            if deleted_keys.contains(key) {
1469                continue;
1470            }
1471            let mut checked: crate::EntityValues = record.clone().into();
1472            checked
1473                .entry("id".to_owned())
1474                .or_insert_with(|| key.id.clone());
1475            checked_changes.insert(key.clone(), checked);
1476        }
1477
1478        // Plan updates and inserts before any mutation. A sparse ledger CREATE
1479        // must be validated against the exact payload sent to SQL, not merely
1480        // against the typed object's default-filled in-memory snapshot.
1481        let mut update_batches: std::collections::BTreeMap<
1482            (String, String),
1483            Vec<crate::EntityKey>,
1484        > = std::collections::BTreeMap::new();
1485        let mut insert_batches: std::collections::BTreeMap<String, Vec<crate::EntityKey>> =
1486            std::collections::BTreeMap::new();
1487
1488        for (key, record) in &checked_changes {
1489            if deleted_keys.contains(key) {
1490                continue;
1491            }
1492            let mut is_new = new_keys.contains(key);
1493
1494            if !is_new {
1495                let descriptor = self
1496                    .data_service
1497                    .metadata
1498                    .context
1499                    .require_entity(&key.entity)
1500                    .map_err(DataServiceError::Runtime)?;
1501                let id_property = descriptor.id_property().ok_or_else(|| {
1502                    DataServiceError::Runtime(RuntimeError::Graph(format!(
1503                        "entity {} has no id property",
1504                        key.entity
1505                    )))
1506                })?;
1507                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1508                let current_row = self
1509                    .fetch_graph_current_row_internal(
1510                        &key.entity,
1511                        &id_property.name,
1512                        &key.id,
1513                        my_trace,
1514                    )
1515                    .await?;
1516                if current_row.is_none() {
1517                    is_new = true;
1518                } else if descriptor.version_property().is_some()
1519                    && root.get_original_version(key).is_none()
1520                {
1521                    return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
1522                        "cannot update {}({:?}) without its loaded original version; load the full entity before mutation",
1523                        key.entity, key.id
1524                    ))));
1525                }
1526            }
1527
1528            match is_new {
1529                true => {
1530                    insert_batches
1531                        .entry(key.entity.to_string())
1532                        .or_default()
1533                        .push(key.clone());
1534                }
1535                false => {
1536                    let mut fields: Vec<String> = record.keys().cloned().collect();
1537                    fields.sort();
1538                    let signature = fields.join(",");
1539                    update_batches
1540                        .entry((key.entity.to_string(), signature))
1541                        .or_default()
1542                        .push(key.clone());
1543                }
1544            }
1545        }
1546
1547        let ordered_insert_keys = self
1548            .order_new_ledger_keys(insert_batches.values().flatten().cloned(), &checked_changes)
1549            .map_err(DataServiceError::Runtime)?;
1550        let mut ordered_insert_batches = Vec::<(String, Vec<crate::EntityKey>)>::new();
1551        for key in ordered_insert_keys {
1552            let entity = key.entity.to_string();
1553            if let Some((batch_entity, keys)) = ordered_insert_batches.last_mut()
1554                && *batch_entity == entity
1555            {
1556                keys.push(key);
1557            } else {
1558                ordered_insert_batches.push((entity, vec![key]));
1559            }
1560        }
1561
1562        let mut prepared_insert_batches = Vec::new();
1563        for (entity, keys) in ordered_insert_batches {
1564            let descriptor = self
1565                .data_service
1566                .metadata
1567                .context
1568                .require_entity(&entity)
1569                .map_err(DataServiceError::Runtime)?;
1570            let mut cmd = teaql_core::BatchInsertCommand::new(&descriptor.name);
1571            let mut traces = Vec::new();
1572            for key in &keys {
1573                let record = checked_changes.get(key).unwrap();
1574                let mut db_record = crate::EntityValues::new();
1575                let mut real_id = key.id.clone();
1576                if crate::data_service::helpers::is_unassigned_id_value(&real_id) {
1577                    let gen_id = self
1578                        .data_service
1579                        .metadata
1580                        .context
1581                        .next_id(&entity)
1582                        .map_err(DataServiceError::Runtime)?;
1583                    real_id = Value::U64(gen_id);
1584                    generated_ids.insert(key.clone(), real_id.clone());
1585                }
1586                db_record.insert("id".to_owned(), real_id);
1587                for (field, value) in record {
1588                    if field == "id" {
1589                        continue;
1590                    }
1591                    db_record.insert(field.clone(), value.clone());
1592                }
1593                crate::data_service::helpers::ensure_initial_version(&mut db_record, descriptor);
1594                crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, true);
1595                let location = locations.get(key).cloned().unwrap_or_default();
1596                self.data_service
1597                    .metadata
1598                    .context
1599                    .validate_required_create_payload(&entity, &db_record, &location)
1600                    .map_err(DataServiceError::Runtime)?;
1601                cmd.batch_values.push(db_record.into());
1602                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1603                traces.push(my_trace);
1604            }
1605            cmd.trace_chains = traces;
1606            prepared_insert_batches.push(cmd);
1607        }
1608
1609        // The required-field gate above is complete before Deletes, Inserts,
1610        // or Updates. Ledger deletion is versioned soft delete, so physical
1611        // foreign keys do not impose child-before-parent ordering here.
1612        for key in &deleted_keys {
1613            if new_keys.contains(key) {
1614                continue;
1615            }
1616            let id = key.id.clone();
1617            let mut cmd = teaql_core::DeleteCommand::new(key.entity.as_ref(), id);
1618            if let Some(version) = root.get_original_version(key) {
1619                cmd = cmd.expected_version(version);
1620            }
1621            cmd.trace_chain = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1622            self.delete_internal(&cmd).await?;
1623        }
1624
1625        for cmd in prepared_insert_batches {
1626            self.execute_prepared_batch_insert(cmd).await?;
1627        }
1628
1629        let mut update_order: Vec<(String, String)> = update_batches.keys().cloned().collect();
1630        update_order.sort();
1631
1632        for signature in update_order {
1633            let keys = update_batches.get(&signature).unwrap();
1634            let descriptor = self
1635                .data_service
1636                .metadata
1637                .context
1638                .require_entity(&signature.0)
1639                .map_err(DataServiceError::Runtime)?;
1640            let mut update_fields: Vec<String> =
1641                signature.1.split(',').map(|s| s.to_string()).collect();
1642            if descriptor
1643                .properties
1644                .iter()
1645                .any(|p| p.name == "update_time")
1646                && !update_fields.contains(&"update_time".to_owned())
1647            {
1648                update_fields.push("update_time".to_owned());
1649            }
1650            if let Some(version_property) = descriptor.version_property()
1651                && !update_fields.contains(&version_property.name)
1652            {
1653                update_fields.push(version_property.name.clone());
1654            }
1655            let mut cmd = teaql_core::BatchUpdateCommand::new(&descriptor.name, update_fields);
1656            let mut traces = Vec::new();
1657            for key in keys {
1658                let record = checked_changes.get(key).unwrap();
1659                let mut db_record = crate::EntityValues::new();
1660                db_record.insert("id".to_owned(), key.id.clone());
1661                for (field, value) in record {
1662                    if field == "id" {
1663                        continue;
1664                    }
1665                    db_record.insert(field.clone(), value.clone());
1666                }
1667                crate::data_service::helpers::increment_version(
1668                    &mut db_record,
1669                    descriptor,
1670                    root.get_original_version(key),
1671                );
1672                crate::data_service::helpers::ensure_timestamps(&mut db_record, descriptor, false);
1673                cmd.batch_values.push(db_record.into());
1674                cmd.batch_ids.push(key.id.clone());
1675                cmd.batch_expected_versions
1676                    .push(root.get_original_version(key));
1677                cmd.batch_old_values.push(None); // or fetch from original state if needed
1678                let my_trace = resolve_trace_chain(root.get_trace_chain(key), &trace_chain);
1679                traces.push(my_trace);
1680            }
1681            cmd.trace_chains = traces;
1682            self.execute_prepared_batch_update(cmd).await?;
1683        }
1684
1685        Ok(generated_ids)
1686    }
1687}