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 if let Some(fields) = &node.dirty_fields {
815 node.values.insert(
816 "_dirty_fields".to_owned(),
817 Value::List(fields.iter().cloned().map(Value::Text).collect()),
818 );
819 }
820 let result = context.check_and_fix_values_at(&node.entity, &mut node.values, location);
821 crate::clear_entity_status(&mut node.values);
822 node.values.remove("_dirty_fields");
823 result?;
824
825 if let Some(root) = root
826 && let Some(id) = node.values.get("id").cloned()
827 {
828 let key = crate::EntityKey::new(node.entity.clone(), id);
829 for (field, value) in &node.values {
830 if before.get(field) != Some(value) {
831 root.set(key.clone(), field.clone(), value.clone());
832 }
833 }
834 }
835 }
836
837 for (relation, children) in &mut node.relations {
838 for (index, child) in children.iter_mut().enumerate() {
839 let child_location = location.clone().member(relation).element(index);
840 preflight_graph(context, child, &child_location, root)?;
841 }
842 }
843 Ok(())
844}
845
846fn ledger_object_locations(node: &GraphNode) -> BTreeMap<crate::EntityKey, ObjectLocation> {
851 fn visit(
852 node: &GraphNode,
853 location: &ObjectLocation,
854 locations: &mut BTreeMap<crate::EntityKey, ObjectLocation>,
855 ) {
856 if let Some(id) = node.values.get("id").cloned() {
857 let key = crate::EntityKey::new(node.entity.clone(), id);
858 locations.entry(key).or_insert_with(|| location.clone());
859 }
860 for (relation, children) in &node.relations {
861 for (index, child) in children.iter().enumerate() {
862 let child_location = location.clone().member(relation).element(index);
863 visit(child, &child_location, locations);
864 }
865 }
866 }
867
868 let mut locations = BTreeMap::new();
869 visit(node, &ObjectLocation::root(), &mut locations);
870 locations
871}
872
873#[cfg(test)]
874mod ledger_location_tests {
875 use super::ledger_object_locations;
876 use crate::{EntityKey, GraphNode};
877 use teaql_core::Value;
878
879 #[test]
880 fn nested_ledger_entity_keeps_model_and_json_error_paths() {
881 let mut order = GraphNode::new("Order");
882 order.values.insert("id".to_owned(), Value::U64(7));
883 let mut line = GraphNode::new("OrderLine");
884 line.values.insert("id".to_owned(), Value::U64(9));
885 order.relations.insert("line_items".to_owned(), vec![line]);
886
887 let locations = ledger_object_locations(&order);
888 assert!(locations[&EntityKey::new("Order", 7_u64)].is_root());
889 let child = &locations[&EntityKey::new("OrderLine", 9_u64)];
890 assert_eq!(child.model_path(), "line_items[0]");
891 assert_eq!(child.instance_path(), "/lineItems/0");
892 }
893}
894
895pub trait AuditedSaveExt {
908 type Entity;
909
910 fn save<'a>(
911 self,
912 context: &'a UserContext,
913 ) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>>;
914}
915
916impl<T> AuditedSaveExt for teaql_core::Audited<T>
917where
918 T: Entity + Send + 'static,
919{
920 type Entity = T;
921
922 fn save<'a>(
923 self,
924 context: &'a UserContext,
925 ) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>> {
926 Box::pin(async move {
927 let _entity_name = T::entity_descriptor().name;
928 let entity = self.into_entity(); let mut node = graph_node_from_entity(context, entity)?;
930 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
931 let saver = context
932 .require_resource::<Arc<dyn DynGraphSaver>>()
933 .map_err(|e| {
934 RuntimeError::Graph(format!(
935 "no DynGraphSaver registered — did you call register_executor()? ({})",
936 e
937 ))
938 })?;
939 let saved = saver.save_graph_dyn(context, node).await?;
940 T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
941 .map_err(|e| RuntimeError::Graph(e.to_string()))
942 })
943 }
944}
945
946#[doc(hidden)]
953pub async fn save_audited_ledger_entity<T>(
954 audited: teaql_core::Audited<T>,
955 context: &UserContext,
956) -> Result<T, RuntimeError>
957where
958 T: crate::LedgerEntity + Send + 'static,
959{
960 let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
961 let result = GRAPH_FIX_TIME
962 .scope(
963 teaql_core::time::Timestamp::now(),
964 GRAPH_FIX_EVIDENCE.scope(
965 evidence.clone(),
966 save_audited_ledger_entity_inner(audited, context),
967 ),
968 )
969 .await;
970 context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
971 result
972}
973
974#[doc(hidden)]
981pub async fn save_audited_ledger_entity_with_executor<T, E>(
982 audited: teaql_core::Audited<T>,
983 context: &UserContext,
984 executor: &E,
985) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
986where
987 T: crate::LedgerEntity + Send + 'static,
988 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
989{
990 let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
991 let result = GRAPH_FIX_TIME
992 .scope(
993 teaql_core::time::Timestamp::now(),
994 GRAPH_FIX_EVIDENCE.scope(
995 evidence.clone(),
996 save_audited_ledger_entity_with_executor_inner(audited, context, executor),
997 ),
998 )
999 .await;
1000 context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
1001 result
1002}
1003
1004async fn save_audited_ledger_entity_with_executor_inner<T, E>(
1005 audited: teaql_core::Audited<T>,
1006 context: &UserContext,
1007 executor: &E,
1008) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
1009where
1010 T: crate::LedgerEntity + Send + 'static,
1011 E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
1012{
1013 let entity = audited.into_entity();
1014 let root = entity.entity_runtime_state();
1015 if let Some(error) = root
1016 .as_ref()
1017 .and_then(|root| root.first_composition_error())
1018 {
1019 return Err(RuntimeError::Graph(format!(
1020 "generated entity graph attachment failed before save: {error}"
1021 )));
1022 }
1023 let mut node = graph_node_from_entity(context, entity)?;
1024
1025 if let Some(root) = root {
1026 let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
1027 let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
1028 reject_cancelled_new_root(&root, &root_key)?;
1029 if let Some(changes) = root.current_change_set().changes().get(&root_key) {
1030 for (field, value) in changes {
1031 node.values.insert(field.clone(), value.clone());
1032 }
1033 }
1034 let mut visited = BTreeSet::from([root_key.clone()]);
1035 hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
1036 preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
1037 merge_relation_mutations_into_root(&root, &node)?;
1038 let has_ledger_changes = !root.current_change_set().changes().is_empty()
1039 || !root.deleted_keys().is_empty()
1040 || !root.new_keys().is_empty();
1041 if has_ledger_changes {
1042 let entity_name = node.entity.clone();
1043 let descriptor = context.require_entity(&entity_name)?;
1044 let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
1045 let id_property = descriptor.id_property().ok_or_else(|| {
1046 RuntimeError::Graph(format!("entity {entity_name} has no id property"))
1047 })?;
1048 let data_service =
1049 crate::EntityDataService::for_executor(context, &entity_name, executor);
1050 let locations = ledger_object_locations(&node);
1051 let generated_ids = data_service
1052 .execute_ledger_plan_internal(root.clone(), &locations)
1053 .await
1054 .map_err(data_service_error_into_runtime)?;
1055
1056 if let Some(new_id) = generated_ids.get(&root_key) {
1057 node.values.insert(id_property.name.clone(), new_id.clone());
1058 }
1059 let persisted_id = node.values.get(&id_property.name).cloned().ok_or_else(|| {
1062 RuntimeError::Graph(format!(
1063 "saved {entity_name} missing identity field {}",
1064 id_property.name
1065 ))
1066 })?;
1067 node.values = data_service
1068 .fetch_graph_current_row_internal(
1069 &entity_name,
1070 &id_property.name,
1071 &persisted_id,
1072 Vec::new(),
1073 )
1074 .await
1075 .map_err(data_service_error_into_runtime)?
1076 .map(Into::into)
1077 .ok_or_else(|| {
1078 RuntimeError::Graph(format!(
1079 "persisted {entity_name} record could not be read back"
1080 ))
1081 })?;
1082 let row = teaql_core::CompactRow::from_map(node.values.into());
1083 let entity = if preserve_relations {
1084 T::from_compact_row_with_context(row, &root)
1085 } else {
1086 T::from_compact_row(row)
1087 }
1088 .map_err(|error| RuntimeError::Graph(error.to_string()))?;
1089 return Ok((entity, Some(root)));
1090 }
1091 }
1092
1093 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
1094 let entity_name = node.entity.clone();
1095 let saved = crate::EntityDataService::for_executor(context, entity_name, executor)
1096 .save_graph_internal(node)
1097 .await
1098 .map_err(data_service_error_into_runtime)?;
1099 let entity = T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
1100 .map_err(|error| RuntimeError::Graph(error.to_string()))?;
1101 Ok((entity, None))
1102}
1103
1104fn data_service_error_into_runtime<E: std::error::Error>(
1105 error: DataServiceError<E>,
1106) -> RuntimeError {
1107 match error {
1108 DataServiceError::Runtime(error) => error,
1109 other => RuntimeError::Graph(other.to_string()),
1110 }
1111}
1112
1113async fn save_audited_ledger_entity_inner<T>(
1114 audited: teaql_core::Audited<T>,
1115 context: &UserContext,
1116) -> Result<T, RuntimeError>
1117where
1118 T: crate::LedgerEntity + Send + 'static,
1119{
1120 let _entity_name = T::entity_descriptor().name;
1121 let entity = audited.into_entity();
1122 let root = entity.entity_runtime_state();
1123 if let Some(error) = root
1124 .as_ref()
1125 .and_then(|root| root.first_composition_error())
1126 {
1127 return Err(RuntimeError::Graph(format!(
1128 "generated entity graph attachment failed before save: {error}"
1129 )));
1130 }
1131 let mut node = graph_node_from_entity(context, entity)?;
1132 let saver = context
1133 .require_resource::<Arc<dyn DynGraphSaver>>()
1134 .map_err(|e| {
1135 RuntimeError::Graph(format!(
1136 "no DynGraphSaver registered — did you call register_executor()? ({e})"
1137 ))
1138 })?;
1139
1140 if let Some(root) = root {
1141 let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
1142 let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
1143 reject_cancelled_new_root(&root, &root_key)?;
1144 if let Some(changes) = root.current_change_set().changes().get(&root_key) {
1145 for (field, value) in changes {
1146 node.values.insert(field.clone(), value.clone());
1147 }
1148 }
1149 let mut visited = BTreeSet::from([root_key.clone()]);
1150 hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
1151 preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
1152 merge_relation_mutations_into_root(&root, &node)?;
1153 let has_ledger_changes = !root.current_change_set().changes().is_empty()
1154 || !root.deleted_keys().is_empty()
1155 || !root.new_keys().is_empty();
1156 if has_ledger_changes {
1157 let descriptor = context.require_entity(&node.entity)?;
1158 let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
1159 let saved = saver.save_ledger_dyn(context, node, root.clone()).await?;
1160 let row = teaql_core::CompactRow::from_map(saved.values.into());
1161 return if preserve_relations {
1162 T::from_compact_row_with_context(row, &root)
1163 } else {
1164 T::from_compact_row(row)
1165 }
1166 .map_err(|e| RuntimeError::Graph(e.to_string()));
1167 }
1168 }
1169
1170 preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
1171 let saved = saver.save_graph_dyn(context, node).await?;
1172 T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
1173 .map_err(|e| RuntimeError::Graph(e.to_string()))
1174}
1175
1176#[cfg(test)]
1177mod tests {
1178 use std::sync::{
1179 Arc,
1180 atomic::{AtomicBool, Ordering},
1181 };
1182
1183 use teaql_core::{Entity, Record, TeaqlEntity, Value};
1184 use teaql_macros::{TeaqlEntity as DeriveTeaqlEntity, teaql_entity};
1185
1186 use super::*;
1187 use crate::{
1188 CheckObjectStatus, CheckResults, InMemoryCheckerRegistry, InMemoryMetadataStore,
1189 TypedChecker, TypedEntityChecker,
1190 };
1191
1192 #[teaql_entity]
1193 #[derive(Clone, Debug, PartialEq, DeriveTeaqlEntity)]
1194 #[teaql(entity = "DirtyTrackedOrder", table = "dirty_tracked_order")]
1195 struct DirtyTrackedOrder {
1196 #[teaql(id)]
1197 id: u64,
1198 #[teaql(version)]
1199 version: i64,
1200 code: String,
1201 }
1202
1203 struct DirtyFieldProbe(Arc<AtomicBool>);
1204
1205 impl TypedChecker<DirtyTrackedOrder> for DirtyFieldProbe {
1206 fn check_and_fix_typed(
1207 &self,
1208 _context: &UserContext,
1209 entity: &mut DirtyTrackedOrder,
1210 status: CheckObjectStatus,
1211 _location: &ObjectLocation,
1212 _results: &mut CheckResults,
1213 ) {
1214 let saw_code = status.is_update()
1215 && entity
1216 .dirty_fields()
1217 .is_some_and(|fields| fields.contains("code"));
1218 self.0.store(saw_code, Ordering::SeqCst);
1219 }
1220 }
1221
1222 #[test]
1223 fn graph_preflight_preserves_dirty_fields_for_typed_checker() {
1224 let observed = Arc::new(AtomicBool::new(false));
1225 let context = UserContext::new()
1226 .with_metadata(
1227 InMemoryMetadataStore::new().with_entity(DirtyTrackedOrder::entity_descriptor()),
1228 )
1229 .with_checker_registry(InMemoryCheckerRegistry::new().with_checker(
1230 TypedEntityChecker::<DirtyTrackedOrder, _>::new(DirtyFieldProbe(observed.clone())),
1231 ));
1232 let row = Record::from([
1233 ("id".to_owned(), Value::U64(7)),
1234 ("version".to_owned(), Value::I64(3)),
1235 ("code".to_owned(), Value::Text("ORD-OLD".to_owned())),
1236 ]);
1237 let entity = DirtyTrackedOrder::from_compact_row(teaql_core::CompactRow::from_map(row))
1238 .expect("materialize tracked entity");
1239 entity.__teaql_runtime_state().set(
1240 crate::EntityKey::new("DirtyTrackedOrder", 7_u64),
1241 "code",
1242 Value::Text("ORD-NEW".to_owned()),
1243 );
1244
1245 let mut node = graph_node_from_entity(&context, entity).expect("build graph node");
1246 preflight_graph(&context, &mut node, &ObjectLocation::root(), None)
1247 .expect("preflight succeeds");
1248
1249 assert!(observed.load(Ordering::SeqCst));
1250 assert!(!node.values.contains_key("_dirty_fields"));
1251 }
1252}