#![allow(clippy::items_after_test_module)]
use std::collections::{BTreeMap, BTreeSet};
use std::future::Future;
use std::marker::PhantomData;
use std::pin::Pin;
use std::sync::Arc;
use teaql_core::{Entity, MutationValues, Value};
use crate::{
DataServiceError, GraphNode, GraphOperation, ObjectLocation, RuntimeError, UserContext,
};
tokio::task_local! {
static GRAPH_FIX_TIME: teaql_core::time::Timestamp;
static GRAPH_FIX_EVIDENCE: Arc<std::sync::Mutex<Vec<crate::FixEvidence>>>;
}
pub(crate) fn current_graph_fix_time() -> teaql_core::time::Timestamp {
GRAPH_FIX_TIME
.try_with(|value| *value)
.unwrap_or_else(|_| teaql_core::time::Timestamp::now())
}
pub(crate) fn record_graph_fix_evidence(evidence: crate::FixEvidence) {
let _ = GRAPH_FIX_EVIDENCE.try_with(|current| current.lock().unwrap().push(evidence));
}
pub(crate) trait DynGraphSaver: Send + Sync {
fn save_graph_dyn<'a>(
&'a self,
context: &'a UserContext,
node: GraphNode,
) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>>;
fn save_ledger_dyn<'a>(
&'a self,
context: &'a UserContext,
node: GraphNode,
root: crate::EntityRuntimeState,
) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>>;
}
pub(crate) struct GraphSaverFor<E> {
_marker: PhantomData<fn() -> E>,
}
impl<E> GraphSaverFor<E> {
pub(crate) fn new() -> Self {
Self {
_marker: PhantomData,
}
}
}
impl<E> DynGraphSaver for GraphSaverFor<E>
where
E: teaql_data_service::QueryExecutor
+ teaql_data_service::MutationExecutor
+ teaql_data_service::TransactionExecutor
+ Send
+ Sync
+ 'static,
for<'tx> <E as teaql_data_service::TransactionExecutor>::Tx<'tx>: Send + Sync,
{
fn save_graph_dyn<'a>(
&'a self,
context: &'a UserContext,
node: GraphNode,
) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>> {
Box::pin(async move {
let entity = node.entity.clone();
let executor = context
.require_resource::<E>()
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
let tx = teaql_data_service::TransactionExecutor::begin(executor)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
let result = {
let eds = crate::EntityDataService::for_executor(context, entity, &tx);
eds.save_graph_internal(node).await
};
match result {
Ok(saved) => {
teaql_data_service::Transaction::commit(tx)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
Ok(saved)
}
Err(error) => {
teaql_data_service::Transaction::rollback(tx)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
Err(match error {
DataServiceError::Runtime(r) => r,
other => RuntimeError::Graph(other.to_string()),
})
}
}
})
}
fn save_ledger_dyn<'a>(
&'a self,
context: &'a UserContext,
mut node: GraphNode,
root: crate::EntityRuntimeState,
) -> Pin<Box<dyn Future<Output = Result<GraphNode, RuntimeError>> + Send + 'a>> {
Box::pin(async move {
let entity = node.entity.clone();
let executor = context
.require_resource::<E>()
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
let descriptor = context.require_entity(&entity)?;
let id_prop = descriptor.id_property().ok_or_else(|| {
RuntimeError::Graph(format!("entity {entity} has no id property"))
})?;
let current_id = node
.values
.get(&id_prop.name)
.cloned()
.unwrap_or(Value::I64(0));
let root_key = crate::EntityKey::new(entity.clone(), current_id);
reject_cancelled_new_root(&root, &root_key)?;
let tx = teaql_data_service::TransactionExecutor::begin(executor)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
let result = async {
let eds = crate::EntityDataService::for_executor(context, &entity, &tx);
let locations = ledger_object_locations(&node);
let generated_ids = eds
.execute_ledger_plan_internal(root.clone(), &locations)
.await?;
if let Some(new_id) = generated_ids.get(&root_key) {
node.values.insert(id_prop.name.clone(), new_id.clone());
}
let persisted_id = node.values.get(&id_prop.name).cloned().ok_or_else(|| {
DataServiceError::Runtime(RuntimeError::Graph(format!(
"saved {entity} missing identity field {}",
id_prop.name
)))
})?;
node.values = eds
.fetch_graph_current_row_internal(
&entity,
&id_prop.name,
&persisted_id,
Vec::new(),
)
.await?
.map(Into::into)
.ok_or_else(|| {
DataServiceError::Runtime(RuntimeError::Graph(format!(
"persisted {entity} record could not be read back"
)))
})?;
Ok(())
}
.await;
match result {
Ok(()) => {
teaql_data_service::Transaction::commit(tx)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
}
Err(error) => {
teaql_data_service::Transaction::rollback(tx)
.await
.map_err(|e| RuntimeError::Graph(e.to_string()))?;
return Err(match error {
DataServiceError::Runtime(r) => r,
other => RuntimeError::Graph(other.to_string()),
});
}
}
root.clear_committed();
Ok(node)
})
}
}
fn reject_cancelled_new_root(
root: &crate::EntityRuntimeState,
root_key: &crate::EntityKey,
) -> Result<(), RuntimeError> {
if root.new_keys().contains(root_key) && root.deleted_keys().contains(root_key) {
return Err(RuntimeError::Graph(format!(
"cancelled new root {root_key:?}: create-then-delete has no persisted entity to return"
)));
}
Ok(())
}
fn can_preserve_loaded_relations(
root: &crate::EntityRuntimeState,
root_key: &crate::EntityKey,
descriptor: &teaql_core::EntityDescriptor,
) -> bool {
if !root.new_keys().is_empty() || !root.deleted_keys().is_empty() {
return false;
}
let changes = root.current_change_set();
changes.changes().iter().all(|(key, fields)| {
key == root_key
&& fields.keys().all(|field| {
!descriptor
.relations
.iter()
.any(|relation| relation.local_key == *field)
})
})
}
#[cfg(test)]
mod save_relation_state_tests {
use super::can_preserve_loaded_relations;
use crate::{EntityKey, EntityRuntimeState};
use teaql_core::{EntityDescriptor, RelationDescriptor, Value};
fn descriptor() -> EntityDescriptor {
let mut descriptor = EntityDescriptor::new("Order");
descriptor
.relations
.push(RelationDescriptor::new("customer", "Customer").local_key("customer_id"));
descriptor
}
#[test]
fn scalar_change_keeps_snapshot_but_relation_changes_invalidate_it() {
let root = EntityRuntimeState::default();
let order = EntityKey::new("Order", Value::I64(1));
let child = EntityKey::new("OrderLine", Value::I64(2));
root.set(order.clone(), "total_amount", Value::I64(100));
assert!(can_preserve_loaded_relations(&root, &order, &descriptor()));
root.set(child, "sku", Value::Text("CHANGED".into()));
assert!(!can_preserve_loaded_relations(&root, &order, &descriptor()));
let root = EntityRuntimeState::default();
root.set(order.clone(), "customer_id", Value::I64(3));
assert!(!can_preserve_loaded_relations(&root, &order, &descriptor()));
}
}
#[cfg(test)]
mod transactional_ledger_readback_tests {
use super::{DynGraphSaver, GraphSaverFor};
use crate::{EntityKey, EntityRuntimeState, GraphNode, InMemoryMetadataStore, UserContext};
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use teaql_core::{DataType, EntityDescriptor, PropertyDescriptor, Value};
use teaql_data_service::{
DataServiceCapabilities, DataServiceExecutor, ExecutionMetadata, MutationExecutor,
MutationRequest, MutationResult, QueryExecutor, QueryRequest, QueryResult, Transaction,
TransactionExecutor,
};
#[derive(Default)]
struct State {
row: BTreeMap<String, Value>,
calls: Vec<&'static str>,
fail_readback: bool,
}
#[derive(Clone)]
struct Ambient(Arc<Mutex<State>>);
struct Tx(Arc<Mutex<State>>);
fn row_result(state: &State) -> QueryResult {
QueryResult {
rows: vec![teaql_core::CompactRow::from_map(state.row.clone())],
metadata: ExecutionMetadata::unrecorded_query(1),
}
}
impl DataServiceExecutor for Ambient {
type Error = std::io::Error;
fn capabilities(&self) -> DataServiceCapabilities {
DataServiceCapabilities {
query: true,
mutation: true,
transaction: true,
..Default::default()
}
}
}
impl DataServiceExecutor for Tx {
type Error = std::io::Error;
fn capabilities(&self) -> DataServiceCapabilities {
Ambient(self.0.clone()).capabilities()
}
}
impl QueryExecutor for Ambient {
async fn query(&self, _request: QueryRequest) -> Result<QueryResult, Self::Error> {
let mut state = self.0.lock().unwrap();
state.calls.push("ambient-query");
Ok(row_result(&state))
}
}
impl QueryExecutor for Tx {
async fn query(&self, _request: QueryRequest) -> Result<QueryResult, Self::Error> {
let mut state = self.0.lock().unwrap();
state.calls.push("transaction-query");
if state.fail_readback && state.calls.contains(&"transaction-mutate") {
return Ok(QueryResult {
rows: Vec::new(),
metadata: ExecutionMetadata::unrecorded_query(0),
});
}
Ok(row_result(&state))
}
}
impl MutationExecutor for Ambient {
async fn mutate(&self, _request: MutationRequest) -> Result<MutationResult, Self::Error> {
panic!("ledger writes must use the transaction executor")
}
}
impl MutationExecutor for Tx {
async fn mutate(&self, request: MutationRequest) -> Result<MutationResult, Self::Error> {
let mut state = self.0.lock().unwrap();
state.calls.push("transaction-mutate");
match request {
MutationRequest::Update(command) => {
assert_eq!(command.expected_version, Some(1));
for (field, value) in command.values {
state.row.insert(field, value);
}
}
MutationRequest::Delete(command) => {
assert_eq!(command.expected_version, Some(1));
state.row.insert("version".to_owned(), Value::I64(-2));
}
other => panic!("unexpected mutation: {other:?}"),
}
Ok(MutationResult {
affected_rows: 1,
generated_values: Default::default(),
persisted_snapshot: None,
metadata: ExecutionMetadata::unrecorded_query(0),
})
}
}
impl TransactionExecutor for Ambient {
type Tx<'a> = Tx;
async fn begin(&self) -> Result<Self::Tx<'_>, Self::Error> {
self.0.lock().unwrap().calls.push("begin");
Ok(Tx(self.0.clone()))
}
}
impl Transaction for Tx {
type Error = std::io::Error;
async fn commit(self) -> Result<(), Self::Error> {
let mut state = self.0.lock().unwrap();
state.calls.push("commit");
state.row.insert("version".to_owned(), Value::I64(3));
state
.row
.insert("name".to_owned(), Value::Text("other writer".to_owned()));
Ok(())
}
async fn rollback(self) -> Result<(), Self::Error> {
let mut state = self.0.lock().unwrap();
state.calls.push("rollback");
state.row.insert("version".to_owned(), Value::I64(1));
state
.row
.insert("name".to_owned(), Value::Text("before".to_owned()));
Ok(())
}
}
#[tokio::test]
async fn ledger_save_returns_transaction_snapshot_before_concurrent_commit_race() {
let state = Arc::new(Mutex::new(State {
row: BTreeMap::from([
("id".to_owned(), Value::I64(1)),
("version".to_owned(), Value::I64(1)),
("name".to_owned(), Value::Text("before".to_owned())),
]),
calls: Vec::new(),
fail_readback: false,
}));
let descriptor = EntityDescriptor::new("Task")
.property(PropertyDescriptor::new("id", DataType::I64).id())
.property(PropertyDescriptor::new("version", DataType::I64).version())
.property(PropertyDescriptor::new("name", DataType::Text));
let context = UserContext::default()
.with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
let root = EntityRuntimeState::default();
let key = EntityKey::new_static("Task", 1_i64);
root.set_original_version(key.clone(), 1);
root.set(key, "name", Value::Text("updated".to_owned()));
let node = GraphNode::new("Task")
.value("id", Value::I64(1))
.value("version", Value::I64(1))
.value("name", Value::Text("before".to_owned()));
let mut context = context;
context.insert_resource(Ambient(state.clone()));
let saved = GraphSaverFor::<Ambient>::new()
.save_ledger_dyn(&context, node, root)
.await
.unwrap();
assert_eq!(saved.values.get("version"), Some(&Value::I64(2)));
assert_eq!(
saved.values.get("name"),
Some(&Value::Text("updated".to_owned()))
);
let state = state.lock().unwrap();
assert_eq!(state.row.get("version"), Some(&Value::I64(3)));
assert_eq!(state.calls.last(), Some(&"commit"));
assert!(!state.calls.contains(&"ambient-query"));
}
#[tokio::test]
async fn failed_authoritative_readback_rolls_back_before_reporting_failure() {
let state = Arc::new(Mutex::new(State {
row: BTreeMap::from([
("id".to_owned(), Value::I64(1)),
("version".to_owned(), Value::I64(1)),
("name".to_owned(), Value::Text("before".to_owned())),
]),
calls: Vec::new(),
fail_readback: true,
}));
let descriptor = EntityDescriptor::new("Task")
.property(PropertyDescriptor::new("id", DataType::I64).id())
.property(PropertyDescriptor::new("version", DataType::I64).version())
.property(PropertyDescriptor::new("name", DataType::Text));
let mut context = UserContext::default()
.with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
context.insert_resource(Ambient(state.clone()));
let root = EntityRuntimeState::default();
let key = EntityKey::new_static("Task", 1_i64);
root.set_original_version(key.clone(), 1);
root.set(key, "name", Value::Text("updated".to_owned()));
let node = GraphNode::new("Task")
.value("id", Value::I64(1))
.value("version", Value::I64(1))
.value("name", Value::Text("before".to_owned()));
let error = GraphSaverFor::<Ambient>::new()
.save_ledger_dyn(&context, node, root.clone())
.await
.unwrap_err();
assert!(error.to_string().contains("could not be read back"));
let state = state.lock().unwrap();
assert_eq!(state.calls.last(), Some(&"rollback"));
assert!(!state.calls.contains(&"commit"));
assert_eq!(state.row.get("version"), Some(&Value::I64(1)));
assert_eq!(
root.get_original_version(&EntityKey::new_static("Task", 1_i64)),
Some(1)
);
}
#[tokio::test]
async fn soft_delete_returns_authoritative_tombstone_before_commit() {
let state = Arc::new(Mutex::new(State {
row: BTreeMap::from([
("id".to_owned(), Value::I64(1)),
("version".to_owned(), Value::I64(1)),
("name".to_owned(), Value::Text("before".to_owned())),
]),
calls: Vec::new(),
fail_readback: false,
}));
let descriptor = EntityDescriptor::new("Task")
.property(PropertyDescriptor::new("id", DataType::I64).id())
.property(PropertyDescriptor::new("version", DataType::I64).version())
.property(PropertyDescriptor::new("name", DataType::Text));
let mut context = UserContext::default()
.with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
context.insert_resource(Ambient(state.clone()));
let root = EntityRuntimeState::default();
let key = EntityKey::new_static("Task", 1_i64);
root.set_original_version(key.clone(), 1);
root.mark_as_delete(key);
let node = GraphNode::new("Task")
.value("id", Value::I64(1))
.value("version", Value::I64(1))
.value("name", Value::Text("before".to_owned()));
let saved = GraphSaverFor::<Ambient>::new()
.save_ledger_dyn(&context, node, root)
.await
.unwrap();
assert_eq!(saved.values.get("version"), Some(&Value::I64(-2)));
let state = state.lock().unwrap();
assert_eq!(state.calls.last(), Some(&"commit"));
assert!(!state.calls.contains(&"ambient-query"));
}
#[tokio::test]
async fn cancelled_new_root_cannot_return_a_fictitious_persisted_entity() {
let state = Arc::new(Mutex::new(State::default()));
let descriptor = EntityDescriptor::new("Task")
.property(PropertyDescriptor::new("id", DataType::I64).id())
.property(PropertyDescriptor::new("version", DataType::I64).version())
.property(PropertyDescriptor::new("name", DataType::Text));
let mut context = UserContext::default()
.with_metadata(InMemoryMetadataStore::new().with_entity(descriptor));
context.insert_resource(Ambient(state.clone()));
let root = EntityRuntimeState::default();
let key = EntityKey::new_static("Task", 7_i64);
root.mark_as_new(key.clone());
root.mark_as_delete(key);
let node = GraphNode::new("Task")
.value("id", Value::I64(7))
.value("version", Value::I64(0))
.value("name", Value::Text("cancelled".to_owned()));
let error = GraphSaverFor::<Ambient>::new()
.save_ledger_dyn(&context, node, root)
.await
.expect_err("cancelled root has no persisted row to return");
assert!(error.to_string().contains("cancelled new root"));
let calls = &state.lock().unwrap().calls;
assert!(!calls.contains(&"transaction-mutate"));
assert!(!calls.contains(&"commit"));
}
}
pub fn graph_node_from_entity<T: Entity>(
context: &UserContext,
entity: T,
) -> Result<GraphNode, RuntimeError> {
let descriptor = T::entity_descriptor();
let loaded_fields = descriptor
.properties
.iter()
.filter(|property| entity.is_field_loaded(&property.name))
.map(|property| Value::Text(property.name.clone()))
.collect::<Vec<_>>();
let dirty_fields = entity.dirty_fields();
let original_values = entity.original_values();
let is_new = entity.is_new();
let is_deleted = entity.is_marked_as_delete();
let comment = entity.get_comment();
let mut node = graph_node_from_values(context, &descriptor.name, entity.into_values())?;
node.values
.insert("_loaded_fields".to_owned(), Value::List(loaded_fields));
node.dirty_fields = dirty_fields;
node.original_values = original_values;
if is_new {
node.operation = GraphOperation::Create;
}
if is_deleted {
node.operation = GraphOperation::Remove;
node.relations.clear();
}
if let Some(c) = comment {
node.set_comment(c);
}
Ok(node)
}
fn graph_node_from_values(
context: &UserContext,
entity: &str,
values: MutationValues,
) -> Result<GraphNode, RuntimeError> {
let descriptor = context.require_entity(entity)?;
let mut node = GraphNode::new(entity);
for (field, value) in values {
if field == "_comment" {
if let Value::Text(comment) = value {
node.set_comment(comment);
}
continue;
}
if field == "_dirty_fields" {
if let Value::List(fields) = value {
let mut dirty = BTreeSet::new();
for f in fields {
if let Value::Text(t) = f {
dirty.insert(t);
}
}
node.dirty_fields = Some(dirty);
}
continue;
}
if field == "_original_values" {
if let Value::Object(orig) = value {
node.original_values = Some(orig.into());
}
continue;
}
if field == "_is_new" {
if matches!(value, Value::Bool(true)) {
node.operation = GraphOperation::Create;
}
continue;
}
if field == "_is_deleted" {
if matches!(value, Value::Bool(true)) {
node.operation = GraphOperation::Remove;
}
continue;
}
let Some(relation) = descriptor.relation_by_name(&field) else {
node.values.insert(field, value);
continue;
};
match value {
Value::Null => {
node.relations.entry(field).or_default();
}
Value::Object(record) => {
let child =
graph_node_from_values(context, &relation.target_entity, record.into())?;
node.relations.entry(field).or_default().push(child);
}
Value::List(values) => {
let children = node.relations.entry(field.clone()).or_default();
for value in values {
let Value::Object(record) = value else {
return Err(RuntimeError::Graph(format!(
"relation {}.{} expects object children, got {:?}",
entity, field, value
)));
};
children.push(graph_node_from_values(
context,
&relation.target_entity,
record.into(),
)?);
}
}
other => {
return Err(RuntimeError::Graph(format!(
"relation {}.{} expects object/list/null, got {:?}",
entity, field, other
)));
}
}
}
Ok(node)
}
fn merge_relation_mutations_into_root(
root: &crate::EntityRuntimeState,
node: &GraphNode,
) -> Result<(), RuntimeError> {
for children in node.relations.values() {
for child in children {
let id = child.values.get("id").cloned().ok_or_else(|| {
RuntimeError::Graph(format!(
"related mutation {} is missing its id",
child.entity
))
})?;
let key = crate::EntityKey::new(child.entity.clone(), id);
match child.operation {
GraphOperation::Create => {
root.mark_as_new(key.clone());
for (field, value) in &child.values {
root.set(key.clone(), field, value.clone());
}
}
GraphOperation::Upsert => {
if let Some(fields) = &child.dirty_fields {
for field in fields {
if let Some(value) = child.values.get(field) {
root.set(key.clone(), field, value.clone());
}
}
}
}
GraphOperation::Remove => root.mark_as_delete(key.clone()),
GraphOperation::Reference => {}
}
if let Some(version) = child
.original_values
.as_ref()
.and_then(|values| values.get("version"))
.and_then(Value::try_i64)
{
root.set_original_version(key, version);
}
merge_relation_mutations_into_root(root, child)?;
}
}
Ok(())
}
fn hydrate_ledger_relations(
context: &UserContext,
root: &crate::EntityRuntimeState,
node: &mut GraphNode,
visited: &mut BTreeSet<crate::EntityKey>,
) -> Result<(), RuntimeError> {
let descriptor = context.require_entity(&node.entity)?;
for relation in &descriptor.relations {
let Some(local_value) = node.values.get(&relation.local_key).cloned() else {
continue;
};
let existing = node.relations.entry(relation.name.clone()).or_default();
let existing_keys = existing
.iter()
.filter_map(|child| {
child
.values
.get("id")
.cloned()
.map(|id| crate::EntityKey::new(child.entity.clone(), id))
})
.collect::<BTreeSet<_>>();
let mut discovered = Vec::new();
for (key, changes) in root.current_change_set().changes() {
if key.entity.as_ref() != relation.target_entity || existing_keys.contains(key) {
continue;
}
let foreign_value = if relation.foreign_key == "id" {
Some(&key.id)
} else {
changes.get(&relation.foreign_key)
};
if foreign_value != Some(&local_value) || !visited.insert(key.clone()) {
continue;
}
let mut values: crate::EntityValues = changes.clone().into();
values
.entry("id".to_owned())
.or_insert_with(|| key.id.clone());
let operation = if root.deleted_keys().contains(key) {
GraphOperation::Remove
} else if root.new_keys().contains(key) || root.get_original_version(key).is_none() {
GraphOperation::Create
} else {
GraphOperation::Upsert
};
let mut child = GraphNode::new(key.entity.to_string());
child.values = values;
child.operation = operation;
hydrate_ledger_relations(context, root, &mut child, visited)?;
discovered.push(child);
}
existing.extend(discovered);
}
Ok(())
}
fn preflight_graph(
context: &UserContext,
node: &mut GraphNode,
location: &ObjectLocation,
root: Option<&crate::EntityRuntimeState>,
) -> Result<(), RuntimeError> {
if !matches!(
node.operation,
GraphOperation::Remove | GraphOperation::Reference
) {
let before = node.values.clone();
let status = match node.operation {
GraphOperation::Create => crate::CheckObjectStatus::Create,
GraphOperation::Upsert => crate::CheckObjectStatus::Update,
GraphOperation::Remove | GraphOperation::Reference => unreachable!(),
};
crate::mark_entity_status(&mut node.values, status);
let result = context.check_and_fix_values_at(&node.entity, &mut node.values, location);
crate::clear_entity_status(&mut node.values);
result?;
if let Some(root) = root
&& let Some(id) = node.values.get("id").cloned()
{
let key = crate::EntityKey::new(node.entity.clone(), id);
for (field, value) in &node.values {
if before.get(field) != Some(value) {
root.set(key.clone(), field.clone(), value.clone());
}
}
}
}
for (relation, children) in &mut node.relations {
for (index, child) in children.iter_mut().enumerate() {
let child_location = location.clone().member(relation).element(index);
preflight_graph(context, child, &child_location, root)?;
}
}
Ok(())
}
fn ledger_object_locations(node: &GraphNode) -> BTreeMap<crate::EntityKey, ObjectLocation> {
fn visit(
node: &GraphNode,
location: &ObjectLocation,
locations: &mut BTreeMap<crate::EntityKey, ObjectLocation>,
) {
if let Some(id) = node.values.get("id").cloned() {
let key = crate::EntityKey::new(node.entity.clone(), id);
locations.entry(key).or_insert_with(|| location.clone());
}
for (relation, children) in &node.relations {
for (index, child) in children.iter().enumerate() {
let child_location = location.clone().member(relation).element(index);
visit(child, &child_location, locations);
}
}
}
let mut locations = BTreeMap::new();
visit(node, &ObjectLocation::root(), &mut locations);
locations
}
#[cfg(test)]
mod ledger_location_tests {
use super::ledger_object_locations;
use crate::{EntityKey, GraphNode};
use teaql_core::Value;
#[test]
fn nested_ledger_entity_keeps_model_and_json_error_paths() {
let mut order = GraphNode::new("Order");
order.values.insert("id".to_owned(), Value::U64(7));
let mut line = GraphNode::new("OrderLine");
line.values.insert("id".to_owned(), Value::U64(9));
order.relations.insert("line_items".to_owned(), vec![line]);
let locations = ledger_object_locations(&order);
assert!(locations[&EntityKey::new("Order", 7_u64)].is_root());
let child = &locations[&EntityKey::new("OrderLine", 9_u64)];
assert_eq!(child.model_path(), "line_items[0]");
assert_eq!(child.instance_path(), "/lineItems/0");
}
}
pub trait AuditedSaveExt {
type Entity;
fn save<'a>(
self,
context: &'a UserContext,
) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>>;
}
impl<T> AuditedSaveExt for teaql_core::Audited<T>
where
T: Entity + Send + 'static,
{
type Entity = T;
fn save<'a>(
self,
context: &'a UserContext,
) -> Pin<Box<dyn Future<Output = Result<Self::Entity, RuntimeError>> + Send + 'a>> {
Box::pin(async move {
let _entity_name = T::entity_descriptor().name;
let entity = self.into_entity(); let mut node = graph_node_from_entity(context, entity)?;
preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
let saver = context
.require_resource::<Arc<dyn DynGraphSaver>>()
.map_err(|e| {
RuntimeError::Graph(format!(
"no DynGraphSaver registered — did you call register_executor()? ({})",
e
))
})?;
let saved = saver.save_graph_dyn(context, node).await?;
T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
.map_err(|e| RuntimeError::Graph(e.to_string()))
})
}
}
#[doc(hidden)]
pub async fn save_audited_ledger_entity<T>(
audited: teaql_core::Audited<T>,
context: &UserContext,
) -> Result<T, RuntimeError>
where
T: crate::LedgerEntity + Send + 'static,
{
let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
let result = GRAPH_FIX_TIME
.scope(
teaql_core::time::Timestamp::now(),
GRAPH_FIX_EVIDENCE.scope(
evidence.clone(),
save_audited_ledger_entity_inner(audited, context),
),
)
.await;
context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
result
}
#[doc(hidden)]
pub async fn save_audited_ledger_entity_with_executor<T, E>(
audited: teaql_core::Audited<T>,
context: &UserContext,
executor: &E,
) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
where
T: crate::LedgerEntity + Send + 'static,
E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
{
let evidence = Arc::new(std::sync::Mutex::new(Vec::new()));
let result = GRAPH_FIX_TIME
.scope(
teaql_core::time::Timestamp::now(),
GRAPH_FIX_EVIDENCE.scope(
evidence.clone(),
save_audited_ledger_entity_with_executor_inner(audited, context, executor),
),
)
.await;
context.replace_last_fix_evidence(evidence.lock().unwrap().clone());
result
}
async fn save_audited_ledger_entity_with_executor_inner<T, E>(
audited: teaql_core::Audited<T>,
context: &UserContext,
executor: &E,
) -> Result<(T, Option<crate::EntityRuntimeState>), RuntimeError>
where
T: crate::LedgerEntity + Send + 'static,
E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Send + Sync,
{
let entity = audited.into_entity();
let root = entity.entity_runtime_state();
if let Some(error) = root
.as_ref()
.and_then(|root| root.first_composition_error())
{
return Err(RuntimeError::Graph(format!(
"generated entity graph attachment failed before save: {error}"
)));
}
let mut node = graph_node_from_entity(context, entity)?;
if let Some(root) = root {
let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
reject_cancelled_new_root(&root, &root_key)?;
if let Some(changes) = root.current_change_set().changes().get(&root_key) {
for (field, value) in changes {
node.values.insert(field.clone(), value.clone());
}
}
let mut visited = BTreeSet::from([root_key.clone()]);
hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
merge_relation_mutations_into_root(&root, &node)?;
let has_ledger_changes = !root.current_change_set().changes().is_empty()
|| !root.deleted_keys().is_empty()
|| !root.new_keys().is_empty();
if has_ledger_changes {
let entity_name = node.entity.clone();
let descriptor = context.require_entity(&entity_name)?;
let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
let id_property = descriptor.id_property().ok_or_else(|| {
RuntimeError::Graph(format!("entity {entity_name} has no id property"))
})?;
let data_service =
crate::EntityDataService::for_executor(context, &entity_name, executor);
let locations = ledger_object_locations(&node);
let generated_ids = data_service
.execute_ledger_plan_internal(root.clone(), &locations)
.await
.map_err(data_service_error_into_runtime)?;
if let Some(new_id) = generated_ids.get(&root_key) {
node.values.insert(id_property.name.clone(), new_id.clone());
}
let persisted_id = node.values.get(&id_property.name).cloned().ok_or_else(|| {
RuntimeError::Graph(format!(
"saved {entity_name} missing identity field {}",
id_property.name
))
})?;
node.values = data_service
.fetch_graph_current_row_internal(
&entity_name,
&id_property.name,
&persisted_id,
Vec::new(),
)
.await
.map_err(data_service_error_into_runtime)?
.map(Into::into)
.ok_or_else(|| {
RuntimeError::Graph(format!(
"persisted {entity_name} record could not be read back"
))
})?;
let row = teaql_core::CompactRow::from_map(node.values.into());
let entity = if preserve_relations {
T::from_compact_row_with_context(row, &root)
} else {
T::from_compact_row(row)
}
.map_err(|error| RuntimeError::Graph(error.to_string()))?;
return Ok((entity, Some(root)));
}
}
preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
let entity_name = node.entity.clone();
let saved = crate::EntityDataService::for_executor(context, entity_name, executor)
.save_graph_internal(node)
.await
.map_err(data_service_error_into_runtime)?;
let entity = T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
.map_err(|error| RuntimeError::Graph(error.to_string()))?;
Ok((entity, None))
}
fn data_service_error_into_runtime<E: std::error::Error>(
error: DataServiceError<E>,
) -> RuntimeError {
match error {
DataServiceError::Runtime(error) => error,
other => RuntimeError::Graph(other.to_string()),
}
}
async fn save_audited_ledger_entity_inner<T>(
audited: teaql_core::Audited<T>,
context: &UserContext,
) -> Result<T, RuntimeError>
where
T: crate::LedgerEntity + Send + 'static,
{
let _entity_name = T::entity_descriptor().name;
let entity = audited.into_entity();
let root = entity.entity_runtime_state();
if let Some(error) = root
.as_ref()
.and_then(|root| root.first_composition_error())
{
return Err(RuntimeError::Graph(format!(
"generated entity graph attachment failed before save: {error}"
)));
}
let mut node = graph_node_from_entity(context, entity)?;
let saver = context
.require_resource::<Arc<dyn DynGraphSaver>>()
.map_err(|e| {
RuntimeError::Graph(format!(
"no DynGraphSaver registered — did you call register_executor()? ({e})"
))
})?;
if let Some(root) = root {
let root_id = node.values.get("id").cloned().unwrap_or(Value::I64(0));
let root_key = crate::EntityKey::new(node.entity.clone(), root_id);
reject_cancelled_new_root(&root, &root_key)?;
if let Some(changes) = root.current_change_set().changes().get(&root_key) {
for (field, value) in changes {
node.values.insert(field.clone(), value.clone());
}
}
let mut visited = BTreeSet::from([root_key.clone()]);
hydrate_ledger_relations(context, &root, &mut node, &mut visited)?;
preflight_graph(context, &mut node, &ObjectLocation::root(), Some(&root))?;
merge_relation_mutations_into_root(&root, &node)?;
let has_ledger_changes = !root.current_change_set().changes().is_empty()
|| !root.deleted_keys().is_empty()
|| !root.new_keys().is_empty();
if has_ledger_changes {
let descriptor = context.require_entity(&node.entity)?;
let preserve_relations = can_preserve_loaded_relations(&root, &root_key, descriptor);
let saved = saver.save_ledger_dyn(context, node, root.clone()).await?;
let row = teaql_core::CompactRow::from_map(saved.values.into());
return if preserve_relations {
T::from_compact_row_with_context(row, &root)
} else {
T::from_compact_row(row)
}
.map_err(|e| RuntimeError::Graph(e.to_string()));
}
}
preflight_graph(context, &mut node, &ObjectLocation::root(), None)?;
let saved = saver.save_graph_dyn(context, node).await?;
T::from_compact_row(teaql_core::CompactRow::from_map(saved.values.into()))
.map_err(|e| RuntimeError::Graph(e.to_string()))
}