1#![allow(clippy::items_after_test_module)] use std::collections::{BTreeMap, BTreeSet};
4use std::future::Future;
5use std::marker::PhantomData;
6use std::pin::Pin;
7use std::sync::Arc;
8
9use teaql_core::{Entity, MutationValues, Value};
10
11use crate::{
12 DataServiceError, GraphNode, GraphOperation, ObjectLocation, RuntimeError, UserContext,
13};
14
15tokio::task_local! {
16 static GRAPH_FIX_TIME: teaql_core::time::Timestamp;
17 static GRAPH_FIX_EVIDENCE: Arc<std::sync::Mutex<Vec<crate::FixEvidence>>>;
18}
19
20pub(crate) fn current_graph_fix_time() -> teaql_core::time::Timestamp {
21 GRAPH_FIX_TIME
22 .try_with(|value| *value)
23 .unwrap_or_else(|_| teaql_core::time::Timestamp::now())
24}
25
26pub(crate) fn record_graph_fix_evidence(evidence: crate::FixEvidence) {
27 let _ = GRAPH_FIX_EVIDENCE.try_with(|current| current.lock().unwrap().push(evidence));
28}
29
30pub(crate) trait DynGraphSaver: Send + Sync {
40 fn save_graph_dyn<'a>(
41 &'a self,
42 context: &'a UserContext,
43 node: GraphNode,
44 ) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>>;
45
46 fn save_ledger_dyn<'a>(
47 &'a self,
48 context: &'a UserContext,
49 node: GraphNode,
50 root: crate::EntityRuntimeState,
51 ) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>>;
52}
53
54pub(crate) struct GraphSaverFor<E> {
60 _marker: PhantomData<fn() -> E>,
61}
62
63impl<E> GraphSaverFor<E> {
64 pub(crate) fn new() -> Self {
65 Self {
66 _marker: PhantomData,
67 }
68 }
69}
70
71impl<E> DynGraphSaver for GraphSaverFor<E>
72where
73 E: teaql_data_service::QueryExecutor
74 + teaql_data_service::MutationExecutor
75 + teaql_data_service::TransactionExecutor
76 + Send
77 + Sync
78 + 'static,
79 for<'tx> <E as teaql_data_service::TransactionExecutor>::Tx<'tx>: Send + Sync,
80{
81 fn save_graph_dyn<'a>(
82 &'a self,
83 context: &'a UserContext,
84 node: GraphNode,
85 ) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>> {
86 Box::pin(async move {
87 let entity = node.entity.clone();
88 let executor = context
89 .require_resource::<E>()
90 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
91 let tx = teaql_data_service::TransactionExecutor::begin(executor)
92 .await
93 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
94 let result = {
95 let eds = crate::EntityDataService::for_executor(context, entity, &tx);
96 eds.save_graph_internal(node).await
97 };
98 match result {
99 Ok(saved) => {
100 teaql_data_service::Transaction::commit(tx)
101 .await
102 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
103 Ok(saved)
104 }
105 Err(error) => {
106 teaql_data_service::Transaction::rollback(tx)
107 .await
108 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
109 Err(match error {
110 DataServiceError::Runtime(r) => r,
111 other => RuntimeError::Graph(other.to_string()),
112 })
113 }
114 }
115 })
116 }
117
118 fn save_ledger_dyn<'a>(
119 &'a self,
120 context: &'a UserContext,
121 mut node: GraphNode,
122 root: crate::EntityRuntimeState,
123 ) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>> {
124 Box::pin(async move {
125 let entity = node.entity.clone();
126 let executor = context
127 .require_resource::<E>()
128 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
129 let descriptor = context.require_entity(&entity)?;
130 let id_prop = descriptor.id_property().ok_or_else(|| {
131 RuntimeError::Graph(format!("entity {entity} has no id property"))
132 })?;
133 let current_id = node
134 .values
135 .get(&id_prop.name)
136 .cloned()
137 .unwrap_or(Value::I64(0));
138 let root_key = crate::EntityKey::new(entity.clone(), current_id);
139 reject_cancelled_new_root(&root, &root_key)?;
140 let tx = teaql_data_service::TransactionExecutor::begin(executor)
141 .await
142 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
143 let result = async {
144 let eds = crate::EntityDataService::for_executor(context, &entity, &tx);
145 let locations = ledger_object_locations(&node);
146 let generated_ids = eds
147 .execute_ledger_plan_internal(root.clone(), &locations)
148 .await?;
149 if let Some(new_id) = generated_ids.get(&root_key) {
150 node.values.insert(id_prop.name.clone(), new_id.clone());
151 }
152 let persisted_id = node.values.get(&id_prop.name).cloned().ok_or_else(|| {
157 DataServiceError::Runtime(RuntimeError::Graph(format!(
158 "saved {entity} missing identity field {}",
159 id_prop.name
160 )))
161 })?;
162 node.values = eds
163 .fetch_graph_current_row_internal(
164 &entity,
165 &id_prop.name,
166 &persisted_id,
167 Vec::new(),
168 )
169 .await?
170 .map(Into::into)
171 .ok_or_else(|| {
172 DataServiceError::Runtime(RuntimeError::Graph(format!(
173 "persisted {entity} record could not be read back"
174 )))
175 })?;
176 Ok(())
177 }
178 .await;
179 match result {
180 Ok(()) => {
181 teaql_data_service::Transaction::commit(tx)
182 .await
183 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
184 }
185 Err(error) => {
186 teaql_data_service::Transaction::rollback(tx)
187 .await
188 .map_err(|e| RuntimeError::Graph(e.to_string()))?;
189 return Err(match error {
190 DataServiceError::Runtime(r) => r,
191 other => RuntimeError::Graph(other.to_string()),
192 });
193 }
194 }
195 root.clear_committed();
196 Ok(node)
197 })
198 }
199}
200
201fn reject_cancelled_new_root(
202 root: &crate::EntityRuntimeState,
203 root_key: &crate::EntityKey,
204) -> Result<(), RuntimeError> {
205 if root.new_keys().contains(root_key) && root.deleted_keys().contains(root_key) {
206 return Err(RuntimeError::Graph(format!(
207 "cancelled new root {root_key:?}: create-then-delete has no persisted entity to return"
208 )));
209 }
210 Ok(())
211}
212
213fn can_preserve_loaded_relations(
217 root: &crate::EntityRuntimeState,
218 root_key: &crate::EntityKey,
219 descriptor: &teaql_core::EntityDescriptor,
220) -> bool {
221 if !root.new_keys().is_empty() || !root.deleted_keys().is_empty() {
222 return false;
223 }
224 let changes = root.current_change_set();
225 changes.changes().iter().all(|(key, fields)| {
226 key == root_key
227 && fields.keys().all(|field| {
228 !descriptor
229 .relations
230 .iter()
231 .any(|relation| relation.local_key == *field)
232 })
233 })
234}
235
236#[cfg(test)]
237mod save_relation_state_tests {
238 use super::can_preserve_loaded_relations;
239 use crate::{EntityKey, EntityRuntimeState};
240 use teaql_core::{EntityDescriptor, RelationDescriptor, Value};
241
242 fn descriptor() -> EntityDescriptor {
243 let mut descriptor = EntityDescriptor::new("Order");
244 descriptor
245 .relations
246 .push(RelationDescriptor::new("customer", "Customer").local_key("customer_id"));
247 descriptor
248 }
249
250 #[test]
251 fn scalar_change_keeps_snapshot_but_relation_changes_invalidate_it() {
252 let root = EntityRuntimeState::default();
253 let order = EntityKey::new("Order", Value::I64(1));
254 let child = EntityKey::new("OrderLine", Value::I64(2));
255 root.set(order.clone(), "total_amount", Value::I64(100));
256 assert!(can_preserve_loaded_relations(&root, &order, &descriptor()));
257 root.set(child, "sku", Value::Text("CHANGED".into()));
258 assert!(!can_preserve_loaded_relations(&root, &order, &descriptor()));
259
260 let root = EntityRuntimeState::default();
261 root.set(order.clone(), "customer_id", Value::I64(3));
262 assert!(!can_preserve_loaded_relations(&root, &order, &descriptor()));
263 }
264}
265
266#[cfg(test)]
267mod transactional_ledger_readback_tests {
268 use super::{DynGraphSaver, GraphSaverFor};
269 use crate::{EntityKey, EntityRuntimeState, GraphNode, InMemoryMetadataStore, UserContext};
270 use std::collections::BTreeMap;
271 use std::sync::{Arc, Mutex};
272 use teaql_core::{DataType, EntityDescriptor, PropertyDescriptor, Value};
273 use teaql_data_service::{
274 DataServiceCapabilities, DataServiceExecutor, ExecutionMetadata, MutationExecutor,
275 MutationRequest, MutationResult, QueryExecutor, QueryRequest, QueryResult, Transaction,
276 TransactionExecutor,
277 };
278
279 #[derive(Default)]
280 struct State {
281 row: BTreeMap<String, Value>,
282 calls: Vec<&'static str>,
283 fail_readback: bool,
284 }
285
286 #[derive(Clone)]
287 struct Ambient(Arc<Mutex<State>>);
288
289 struct Tx(Arc<Mutex<State>>);
290
291 fn row_result(state: &State) -> QueryResult {
292 QueryResult {
293 rows: vec![teaql_core::CompactRow::from_map(state.row.clone())],
294 metadata: ExecutionMetadata::unrecorded_query(1),
295 }
296 }
297
298 impl DataServiceExecutor for Ambient {
299 type Error = std::io::Error;
300
301 fn capabilities(&self) -> DataServiceCapabilities {
302 DataServiceCapabilities {
303 query: true,
304 mutation: true,
305 transaction: true,
306 ..Default::default()
307 }
308 }
309 }
310
311 impl DataServiceExecutor for Tx {
312 type Error = std::io::Error;
313
314 fn capabilities(&self) -> DataServiceCapabilities {
315 Ambient(self.0.clone()).capabilities()
316 }
317 }
318
319 impl QueryExecutor for Ambient {
320 async fn query(&self, _request: QueryRequest) -> Result<QueryResult, Self::Error> {
321 let mut state = self.0.lock().unwrap();
322 state.calls.push("ambient-query");
323 Ok(row_result(&state))
324 }
325 }
326
327 impl QueryExecutor for Tx {
328 async fn query(&self, _request: QueryRequest) -> Result<QueryResult, Self::Error> {
329 let mut state = self.0.lock().unwrap();
330 state.calls.push("transaction-query");
331 if state.fail_readback && state.calls.contains(&"transaction-mutate") {
332 return Ok(QueryResult {
333 rows: Vec::new(),
334 metadata: ExecutionMetadata::unrecorded_query(0),
335 });
336 }
337 Ok(row_result(&state))
338 }
339 }
340
341 impl MutationExecutor for Ambient {
342 async fn mutate(&self, _request: MutationRequest) -> Result<MutationResult, Self::Error> {
343 panic!("ledger writes must use the transaction executor")
344 }
345 }
346
347 impl MutationExecutor for Tx {
348 async fn mutate(&self, request: MutationRequest) -> Result<MutationResult, Self::Error> {
349 let mut state = self.0.lock().unwrap();
350 state.calls.push("transaction-mutate");
351 match request {
352 MutationRequest::Update(command) => {
353 assert_eq!(command.expected_version, Some(1));
354 for (field, value) in command.values {
355 state.row.insert(field, value);
356 }
357 }
358 MutationRequest::Delete(command) => {
359 assert_eq!(command.expected_version, Some(1));
360 state.row.insert("version".to_owned(), Value::I64(-2));
361 }
362 other => panic!("unexpected mutation: {other:?}"),
363 }
364 Ok(MutationResult {
365 affected_rows: 1,
366 generated_values: Default::default(),
367 persisted_snapshot: None,
368 metadata: ExecutionMetadata::unrecorded_query(0),
369 })
370 }
371 }
372
373 impl TransactionExecutor for Ambient {
374 type Tx<'a> = Tx;
375
376 async fn begin(&self) -> Result<Self::Tx<'_>, Self::Error> {
377 self.0.lock().unwrap().calls.push("begin");
378 Ok(Tx(self.0.clone()))
379 }
380 }
381
382 impl Transaction for Tx {
383 type Error = std::io::Error;
384
385 async fn commit(self) -> Result<(), Self::Error> {
386 let mut state = self.0.lock().unwrap();
387 state.calls.push("commit");
388 state.row.insert("version".to_owned(), Value::I64(3));
390 state
391 .row
392 .insert("name".to_owned(), Value::Text("other writer".to_owned()));
393 Ok(())
394 }
395
396 async fn rollback(self) -> Result<(), Self::Error> {
397 let mut state = self.0.lock().unwrap();
398 state.calls.push("rollback");
399 state.row.insert("version".to_owned(), Value::I64(1));
400 state
401 .row
402 .insert("name".to_owned(), Value::Text("before".to_owned()));
403 Ok(())
404 }
405 }
406
407 #[tokio::test]
408 async fn ledger_save_returns_transaction_snapshot_before_concurrent_commit_race() {
409 let state = Arc::new(Mutex::new(State {
410 row: BTreeMap::from([
411 ("id".to_owned(), Value::I64(1)),
412 ("version".to_owned(), Value::I64(1)),
413 ("name".to_owned(), Value::Text("before".to_owned())),
414 ]),
415 calls: Vec::new(),
416 fail_readback: false,
417 }));
418 let descriptor = EntityDescriptor::new("Task")
419 .property(PropertyDescriptor::new("id", DataType::I64).id())
420 .property(PropertyDescriptor::new("version", DataType::I64).version())
421 .property(PropertyDescriptor::new("name", DataType::Text));
422 let context = UserContext::default()
423 .with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
424 let root = EntityRuntimeState::default();
425 let key = EntityKey::new_static("Task", 1_i64);
426 root.set_original_version(key.clone(), 1);
427 root.set(key, "name", Value::Text("updated".to_owned()));
428 let node = GraphNode::new("Task")
429 .value("id", Value::I64(1))
430 .value("version", Value::I64(1))
431 .value("name", Value::Text("before".to_owned()));
432 let mut context = context;
433 context.insert_resource(Ambient(state.clone()));
434
435 let saved = GraphSaverFor::<Ambient>::new()
436 .save_ledger_dyn(&context, node, root)
437 .await
438 .unwrap();
439 assert_eq!(saved.values.get("version"), Some(&Value::I64(2)));
440 assert_eq!(
441 saved.values.get("name"),
442 Some(&Value::Text("updated".to_owned()))
443 );
444 let state = state.lock().unwrap();
445 assert_eq!(state.row.get("version"), Some(&Value::I64(3)));
446 assert_eq!(state.calls.last(), Some(&"commit"));
447 assert!(!state.calls.contains(&"ambient-query"));
448 }
449
450 #[tokio::test]
451 async fn failed_authoritative_readback_rolls_back_before_reporting_failure() {
452 let state = Arc::new(Mutex::new(State {
453 row: BTreeMap::from([
454 ("id".to_owned(), Value::I64(1)),
455 ("version".to_owned(), Value::I64(1)),
456 ("name".to_owned(), Value::Text("before".to_owned())),
457 ]),
458 calls: Vec::new(),
459 fail_readback: true,
460 }));
461 let descriptor = EntityDescriptor::new("Task")
462 .property(PropertyDescriptor::new("id", DataType::I64).id())
463 .property(PropertyDescriptor::new("version", DataType::I64).version())
464 .property(PropertyDescriptor::new("name", DataType::Text));
465 let mut context = UserContext::default()
466 .with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
467 context.insert_resource(Ambient(state.clone()));
468 let root = EntityRuntimeState::default();
469 let key = EntityKey::new_static("Task", 1_i64);
470 root.set_original_version(key.clone(), 1);
471 root.set(key, "name", Value::Text("updated".to_owned()));
472 let node = GraphNode::new("Task")
473 .value("id", Value::I64(1))
474 .value("version", Value::I64(1))
475 .value("name", Value::Text("before".to_owned()));
476
477 let error = GraphSaverFor::<Ambient>::new()
478 .save_ledger_dyn(&context, node, root.clone())
479 .await
480 .unwrap_err();
481 assert!(error.to_string().contains("could not be read back"));
482 let state = state.lock().unwrap();
483 assert_eq!(state.calls.last(), Some(&"rollback"));
484 assert!(!state.calls.contains(&"commit"));
485 assert_eq!(state.row.get("version"), Some(&Value::I64(1)));
486 assert_eq!(
487 root.get_original_version(&EntityKey::new_static("Task", 1_i64)),
488 Some(1)
489 );
490 }
491
492 #[tokio::test]
493 async fn soft_delete_returns_authoritative_tombstone_before_commit() {
494 let state = Arc::new(Mutex::new(State {
495 row: BTreeMap::from([
496 ("id".to_owned(), Value::I64(1)),
497 ("version".to_owned(), Value::I64(1)),
498 ("name".to_owned(), Value::Text("before".to_owned())),
499 ]),
500 calls: Vec::new(),
501 fail_readback: false,
502 }));
503 let descriptor = EntityDescriptor::new("Task")
504 .property(PropertyDescriptor::new("id", DataType::I64).id())
505 .property(PropertyDescriptor::new("version", DataType::I64).version())
506 .property(PropertyDescriptor::new("name", DataType::Text));
507 let mut context = UserContext::default()
508 .with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
509 context.insert_resource(Ambient(state.clone()));
510 let root = EntityRuntimeState::default();
511 let key = EntityKey::new_static("Task", 1_i64);
512 root.set_original_version(key.clone(), 1);
513 root.mark_as_delete(key);
514 let node = GraphNode::new("Task")
515 .value("id", Value::I64(1))
516 .value("version", Value::I64(1))
517 .value("name", Value::Text("before".to_owned()));
518
519 let saved = GraphSaverFor::<Ambient>::new()
520 .save_ledger_dyn(&context, node, root)
521 .await
522 .unwrap();
523 assert_eq!(saved.values.get("version"), Some(&Value::I64(-2)));
524 let state = state.lock().unwrap();
525 assert_eq!(state.calls.last(), Some(&"commit"));
526 assert!(!state.calls.contains(&"ambient-query"));
527 }
528
529 #[tokio::test]
530 async fn cancelled_new_root_cannot_return_a_fictitious_persisted_entity() {
531 let state = Arc::new(Mutex::new(State::default()));
532 let descriptor = EntityDescriptor::new("Task")
533 .property(PropertyDescriptor::new("id", DataType::I64).id())
534 .property(PropertyDescriptor::new("version", DataType::I64).version())
535 .property(PropertyDescriptor::new("name", DataType::Text));
536 let mut context = UserContext::default()
537 .with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
538 context.insert_resource(Ambient(state.clone()));
539 let root = EntityRuntimeState::default();
540 let key = EntityKey::new_static("Task", 7_i64);
541 root.mark_as_new(key.clone());
542 root.mark_as_delete(key);
543 let node = GraphNode::new("Task")
544 .value("id", Value::I64(7))
545 .value("version", Value::I64(0))
546 .value("name", Value::Text("cancelled".to_owned()));
547
548 let error = GraphSaverFor::<Ambient>::new()
549 .save_ledger_dyn(&context, node, root)
550 .await
551 .expect_err("cancelled root has no persisted row to return");
552 assert!(error.to_string().contains("cancelled new root"));
553 let calls = &state.lock().unwrap().calls;
554 assert!(!calls.contains(&"transaction-mutate"));
555 assert!(!calls.contains(&"commit"));
556 }
557}
558
559pub fn graph_node_from_entity<T: Entity>(
569 context: &UserContext,
570 entity: T,
571) -> Result<GraphNode, RuntimeError> {
572 let descriptor = T::entity_descriptor();
573 let loaded_fields = descriptor
574 .properties
575 .iter()
576 .filter(|property| entity.is_field_loaded(&property.name))
577 .map(|property| Value::Text(property.name.clone()))
578 .collect::<Vec<_>>();
579 let dirty_fields = entity.dirty_fields();
580 let original_values = entity.original_values();
581 let is_new = entity.is_new();
582 let is_deleted = entity.is_marked_as_delete();
583 let comment = entity.get_comment();
584 let mut node = graph_node_from_values(context, &descriptor.name, entity.into_values())?;
585 node.values
586 .insert("_loaded_fields".to_owned(), Value::List(loaded_fields));
587 node.dirty_fields = dirty_fields;
588 node.original_values = original_values;
589 if is_new {
590 node.operation = GraphOperation::Create;
591 }
592 if is_deleted {
593 node.operation = GraphOperation::Remove;
594 node.relations.clear();
595 }
596 if let Some(c) = comment {
597 node.set_comment(c);
598 }
599 Ok(node)
600}
601
602fn graph_node_from_values(
606 context: &UserContext,
607 entity: &str,
608 values: MutationValues,
609) -> Result<GraphNode, RuntimeError> {
610 let descriptor = context.require_entity(entity)?;
611 let mut node = GraphNode::new(entity);
612
613 for (field, value) in values {
614 if field == "_comment" {
615 if let Value::Text(comment) = value {
616 node.set_comment(comment);
617 }
618 continue;
619 }
620 if field == "_dirty_fields" {
621 if let Value::List(fields) = value {
622 let mut dirty = BTreeSet::new();
623 for f in fields {
624 if let Value::Text(t) = f {
625 dirty.insert(t);
626 }
627 }
628 node.dirty_fields = Some(dirty);
629 }
630 continue;
631 }
632 if field == "_original_values" {
633 if let Value::Object(orig) = value {
634 node.original_values = Some(orig.into());
635 }
636 continue;
637 }
638 if field == "_is_new" {
639 if matches!(value, Value::Bool(true)) {
640 node.operation = GraphOperation::Create;
641 }
642 continue;
643 }
644 if field == "_is_deleted" {
645 if matches!(value, Value::Bool(true)) {
646 node.operation = GraphOperation::Remove;
647 }
648 continue;
649 }
650 let Some(relation) = descriptor.relation_by_name(&field) else {
651 node.values.insert(field, value);
652 continue;
653 };
654
655 match value {
656 Value::Null => {
657 node.relations.entry(field).or_default();
658 }
659 Value::Object(record) => {
660 let child =
661 graph_node_from_values(context, &relation.target_entity, record.into())?;
662 node.relations.entry(field).or_default().push(child);
663 }
664 Value::List(values) => {
665 let children = node.relations.entry(field.clone()).or_default();
666 for value in values {
667 let Value::Object(record) = value else {
668 return Err(RuntimeError::Graph(format!(
669 "relation {}.{} expects object children, got {:?}",
670 entity, field, value
671 )));
672 };
673 children.push(graph_node_from_values(
674 context,
675 &relation.target_entity,
676 record.into(),
677 )?);
678 }
679 }
680 other => {
681 return Err(RuntimeError::Graph(format!(
682 "relation {}.{} expects object/list/null, got {:?}",
683 entity, field, other
684 )));
685 }
686 }
687 }
688
689 Ok(node)
690}
691
692fn merge_relation_mutations_into_root(
693 root: &crate::EntityRuntimeState,
694 node: &GraphNode,
695) -> Result<(), RuntimeError> {
696 for children in node.relations.values() {
697 for child in children {
698 let id = child.values.get("id").cloned().ok_or_else(|| {
699 RuntimeError::Graph(format!(
700 "related mutation {} is missing its id",
701 child.entity
702 ))
703 })?;
704 let key = crate::EntityKey::new(child.entity.clone(), id);
705
706 match child.operation {
707 GraphOperation::Create => {
708 root.mark_as_new(key.clone());
709 for (field, value) in &child.values {
710 root.set(key.clone(), field, value.clone());
711 }
712 }
713 GraphOperation::Upsert => {
714 if let Some(fields) = &child.dirty_fields {
715 for field in fields {
716 if let Some(value) = child.values.get(field) {
717 root.set(key.clone(), field, value.clone());
718 }
719 }
720 }
721 }
722 GraphOperation::Remove => root.mark_as_delete(key.clone()),
723 GraphOperation::Reference => {}
724 }
725
726 if let Some(version) = child
727 .original_values
728 .as_ref()
729 .and_then(|values| values.get("version"))
730 .and_then(Value::try_i64)
731 {
732 root.set_original_version(key, version);
733 }
734 merge_relation_mutations_into_root(root, child)?;
735 }
736 }
737 Ok(())
738}
739
740fn hydrate_ledger_relations(
741 context: &UserContext,
742 root: &crate::EntityRuntimeState,
743 node: &mut GraphNode,
744 visited: &mut BTreeSet<crate::EntityKey>,
745) -> Result<(), RuntimeError> {
746 let descriptor = context.require_entity(&node.entity)?;
747 for relation in &descriptor.relations {
748 let Some(local_value) = node.values.get(&relation.local_key).cloned() else {
749 continue;
750 };
751 let existing = node.relations.entry(relation.name.clone()).or_default();
752 let existing_keys = existing
753 .iter()
754 .filter_map(|child| {
755 child
756 .values
757 .get("id")
758 .cloned()
759 .map(|id| crate::EntityKey::new(child.entity.clone(), id))
760 })
761 .collect::<BTreeSet<_>>();
762 let mut discovered = Vec::new();
763 for (key, changes) in root.current_change_set().changes() {
764 if key.entity.as_ref() != relation.target_entity || existing_keys.contains(key) {
765 continue;
766 }
767 let foreign_value = if relation.foreign_key == "id" {
768 Some(&key.id)
769 } else {
770 changes.get(&relation.foreign_key)
771 };
772 if foreign_value != Some(&local_value) || !visited.insert(key.clone()) {
773 continue;
774 }
775 let mut values: crate::EntityValues = changes.clone().into();
776 values
777 .entry("id".to_owned())
778 .or_insert_with(|| key.id.clone());
779 let operation = if root.deleted_keys().contains(key) {
780 GraphOperation::Remove
781 } else if root.new_keys().contains(key) || root.get_original_version(key).is_none() {
782 GraphOperation::Create
783 } else {
784 GraphOperation::Upsert
785 };
786 let mut child = GraphNode::new(key.entity.to_string());
787 child.values = values;
788 child.operation = operation;
789 hydrate_ledger_relations(context, root, &mut child, visited)?;
790 discovered.push(child);
791 }
792 existing.extend(discovered);
793 }
794 Ok(())
795}
796
797fn preflight_graph(
798 context: &UserContext,
799 node: &mut GraphNode,
800 location: &ObjectLocation,
801 root: Option<&crate::EntityRuntimeState>,
802) -> Result<(), RuntimeError> {
803 if !matches!(
804 node.operation,
805 GraphOperation::Remove | GraphOperation::Reference
806 ) {
807 let before = node.values.clone();
808 let status = match node.operation {
809 GraphOperation::Create => crate::CheckObjectStatus::Create,
810 GraphOperation::Upsert => crate::CheckObjectStatus::Update,
811 GraphOperation::Remove | GraphOperation::Reference => unreachable!(),
812 };
813 crate::mark_entity_status(&mut node.values, status);
814 let result = context.check_and_fix_values_at(&node.entity, &mut node.values, location);
815 crate::clear_entity_status(&mut node.values);
816 result?;
817
818 if let Some(root) = root
819 && let Some(id) = node.values.get("id").cloned()
820 {
821 let key = crate::EntityKey::new(node.entity.clone(), id);
822 for (field, value) in &node.values {
823 if before.get(field) != Some(value) {
824 root.set(key.clone(), field.clone(), value.clone());
825 }
826 }
827 }
828 }
829
830 for (relation, children) in &mut node.relations {
831 for (index, child) in children.iter_mut().enumerate() {
832 let child_location = location.clone().member(relation).element(index);
833 preflight_graph(context, child, &child_location, root)?;
834 }
835 }
836 Ok(())
837}
838
839fn ledger_object_locations(node: &GraphNode) -> BTreeMap<crate::EntityKey, ObjectLocation> {
844 fn visit(
845 node: &GraphNode,
846 location: &ObjectLocation,
847 locations: &mut BTreeMap<crate::EntityKey, ObjectLocation>,
848 ) {
849 if let Some(id) = node.values.get("id").cloned() {
850 let key = crate::EntityKey::new(node.entity.clone(), id);
851 locations.entry(key).or_insert_with(|| location.clone());
852 }
853 for (relation, children) in &node.relations {
854 for (index, child) in children.iter().enumerate() {
855 let child_location = location.clone().member(relation).element(index);
856 visit(child, &child_location, locations);
857 }
858 }
859 }
860
861 let mut locations = BTreeMap::new();
862 visit(node, &ObjectLocation::root(), &mut locations);
863 locations
864}
865
866#[cfg(test)]
867mod ledger_location_tests {
868 use super::ledger_object_locations;
869 use crate::{EntityKey, GraphNode};
870 use teaql_core::Value;
871
872 #[test]
873 fn nested_ledger_entity_keeps_model_and_json_error_paths() {
874 let mut order = GraphNode::new("Order");
875 order.values.insert("id".to_owned(), Value::U64(7));
876 let mut line = GraphNode::new("OrderLine");
877 line.values.insert("id".to_owned(), Value::U64(9));
878 order.relations.insert("line_items".to_owned(), vec![line]);
879
880 let locations = ledger_object_locations(&order);
881 assert!(locations[&EntityKey::new("Order", 7_u64)].is_root());
882 let child = &locations[&EntityKey::new("OrderLine", 9_u64)];
883 assert_eq!(child.model_path(), "line_items[0]");
884 assert_eq!(child.instance_path(), "/lineItems/0");
885 }
886}
887
888pub trait AuditedSaveExt {
901 type Entity;
902
903 fn save<'a>(
904 self,
905 context: &'a UserContext,
906 ) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>>;
907}
908
909impl<T> AuditedSaveExt for teaql_core::Audited<T>
910where
911 T: Entity + Send + 'static,
912{
913 type Entity = T;
914
915 fn save<'a>(
916 self,
917 context: &'a UserContext,
918 ) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>> {
919 Box::pin(async move {
920 let _entity_name = T::entity_descriptor().name;
921 let entity = self.into_entity(); let mut node = graph_node_from_entity(context, entity)?;
923 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
924 let saver = context
925 .require_resource::<Arc<dyn DynGraphSaver>>()
926 .map_err(|e| {
927 RuntimeError::Graph(format!(
928 "no DynGraphSaver registered — did you call register_executor()? ({})",
929 e
930 ))
931 })?;
932 let saved = saver.save_graph_dyn(context, node).await?;
933 T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
934 .map_err(|e| RuntimeError::Graph(e.to_string()))
935 })
936 }
937}
938
939#[doc(hidden)]
946pub async fn save_audited_ledger_entity<T>(
947 audited: teaql_core::Audited<T>,
948 context: &UserContext,
949) -> Result<T, RuntimeError>
950where
951 T: crate::LedgerEntity + Send + 'static,
952{
953 let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
954 let result = GRAPH_FIX_TIME
955 .scope(
956 teaql_core::time::Timestamp::now(),
957 GRAPH_FIX_EVIDENCE.scope(
958 evidence.clone(),
959 save_audited_ledger_entity_inner(audited, context),
960 ),
961 )
962 .await;
963 context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
964 result
965}
966
967#[doc(hidden)]
974pub async fn save_audited_ledger_entity_with_executor<T, E>(
975 audited: teaql_core::Audited<T>,
976 context: &UserContext,
977 executor: &E,
978) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
979where
980 T: crate::LedgerEntity + Send + 'static,
981 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
982{
983 let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
984 let result = GRAPH_FIX_TIME
985 .scope(
986 teaql_core::time::Timestamp::now(),
987 GRAPH_FIX_EVIDENCE.scope(
988 evidence.clone(),
989 save_audited_ledger_entity_with_executor_inner(audited, context, executor),
990 ),
991 )
992 .await;
993 context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
994 result
995}
996
997async fn save_audited_ledger_entity_with_executor_inner<T, E>(
998 audited: teaql_core::Audited<T>,
999 context: &UserContext,
1000 executor: &E,
1001) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
1002where
1003 T: crate::LedgerEntity + Send + 'static,
1004 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
1005{
1006 let entity = audited.into_entity();
1007 let root = entity.entity_runtime_state();
1008 if let Some(error) = root
1009 .as_ref()
1010 .and_then(|root| root.first_composition_error())
1011 {
1012 return Err(RuntimeError::Graph(format!(
1013 "generated entity graph attachment failed before save: {error}"
1014 )));
1015 }
1016 let mut node = graph_node_from_entity(context, entity)?;
1017
1018 if let Some(root) = root {
1019 let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
1020 let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
1021 reject_cancelled_new_root(&root, &root_key)?;
1022 if let Some(changes) = root.current_change_set().changes().get(&root_key) {
1023 for (field, value) in changes {
1024 node.values.insert(field.clone(), value.clone());
1025 }
1026 }
1027 let mut visited = BTreeSet::from([root_key.clone()]);
1028 hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
1029 preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
1030 merge_relation_mutations_into_root(&root, &node)?;
1031 let has_ledger_changes = !root.current_change_set().changes().is_empty()
1032 || !root.deleted_keys().is_empty()
1033 || !root.new_keys().is_empty();
1034 if has_ledger_changes {
1035 let entity_name = node.entity.clone();
1036 let descriptor = context.require_entity(&entity_name)?;
1037 let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
1038 let id_property = descriptor.id_property().ok_or_else(|| {
1039 RuntimeError::Graph(format!("entity {entity_name} has no id property"))
1040 })?;
1041 let data_service =
1042 crate::EntityDataService::for_executor(context, &entity_name, executor);
1043 let locations = ledger_object_locations(&node);
1044 let generated_ids = data_service
1045 .execute_ledger_plan_internal(root.clone(), &locations)
1046 .await
1047 .map_err(data_service_error_into_runtime)?;
1048
1049 if let Some(new_id) = generated_ids.get(&root_key) {
1050 node.values.insert(id_property.name.clone(), new_id.clone());
1051 }
1052 let persisted_id = node.values.get(&id_property.name).cloned().ok_or_else(|| {
1055 RuntimeError::Graph(format!(
1056 "saved {entity_name} missing identity field {}",
1057 id_property.name
1058 ))
1059 })?;
1060 node.values = data_service
1061 .fetch_graph_current_row_internal(
1062 &entity_name,
1063 &id_property.name,
1064 &persisted_id,
1065 Vec::new(),
1066 )
1067 .await
1068 .map_err(data_service_error_into_runtime)?
1069 .map(Into::into)
1070 .ok_or_else(|| {
1071 RuntimeError::Graph(format!(
1072 "persisted {entity_name} record could not be read back"
1073 ))
1074 })?;
1075 let row = teaql_core::CompactRow::from_map(node.values.into());
1076 let entity = if preserve_relations {
1077 T::from_compact_row_with_context(row, &root)
1078 } else {
1079 T::from_compact_row(row)
1080 }
1081 .map_err(|error| RuntimeError::Graph(error.to_string()))?;
1082 return Ok((entity, Some(root)));
1083 }
1084 }
1085
1086 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
1087 let entity_name = node.entity.clone();
1088 let saved = crate::EntityDataService::for_executor(context, entity_name, executor)
1089 .save_graph_internal(node)
1090 .await
1091 .map_err(data_service_error_into_runtime)?;
1092 let entity = T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
1093 .map_err(|error| RuntimeError::Graph(error.to_string()))?;
1094 Ok((entity, None))
1095}
1096
1097fn data_service_error_into_runtime<E: std::error::Error>(
1098 error: DataServiceError<E>,
1099) -> RuntimeError {
1100 match error {
1101 DataServiceError::Runtime(error) => error,
1102 other => RuntimeError::Graph(other.to_string()),
1103 }
1104}
1105
1106async fn save_audited_ledger_entity_inner<T>(
1107 audited: teaql_core::Audited<T>,
1108 context: &UserContext,
1109) -> Result<T, RuntimeError>
1110where
1111 T: crate::LedgerEntity + Send + 'static,
1112{
1113 let _entity_name = T::entity_descriptor().name;
1114 let entity = audited.into_entity();
1115 let root = entity.entity_runtime_state();
1116 if let Some(error) = root
1117 .as_ref()
1118 .and_then(|root| root.first_composition_error())
1119 {
1120 return Err(RuntimeError::Graph(format!(
1121 "generated entity graph attachment failed before save: {error}"
1122 )));
1123 }
1124 let mut node = graph_node_from_entity(context, entity)?;
1125 let saver = context
1126 .require_resource::<Arc<dyn DynGraphSaver>>()
1127 .map_err(|e| {
1128 RuntimeError::Graph(format!(
1129 "no DynGraphSaver registered — did you call register_executor()? ({e})"
1130 ))
1131 })?;
1132
1133 if let Some(root) = root {
1134 let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
1135 let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
1136 reject_cancelled_new_root(&root, &root_key)?;
1137 if let Some(changes) = root.current_change_set().changes().get(&root_key) {
1138 for (field, value) in changes {
1139 node.values.insert(field.clone(), value.clone());
1140 }
1141 }
1142 let mut visited = BTreeSet::from([root_key.clone()]);
1143 hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
1144 preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
1145 merge_relation_mutations_into_root(&root, &node)?;
1146 let has_ledger_changes = !root.current_change_set().changes().is_empty()
1147 || !root.deleted_keys().is_empty()
1148 || !root.new_keys().is_empty();
1149 if has_ledger_changes {
1150 let descriptor = context.require_entity(&node.entity)?;
1151 let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
1152 let saved = saver.save_ledger_dyn(context, node, root.clone()).await?;
1153 let row = teaql_core::CompactRow::from_map(saved.values.into());
1154 return if preserve_relations {
1155 T::from_compact_row_with_context(row, &root)
1156 } else {
1157 T::from_compact_row(row)
1158 }
1159 .map_err(|e| RuntimeError::Graph(e.to_string()));
1160 }
1161 }
1162
1163 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
1164 let saved = saver.save_graph_dyn(context, node).await?;
1165 T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
1166 .map_err(|e| RuntimeError::Graph(e.to_string()))
1167}