Skip to main content

teaql_runtime/data_service/
graph.rs

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