#![allow(dead_code)]
use teaql_core::{
CompactRow, DeleteCommand, Entity, InsertCommand, MutationValues, RecoverCommand, SelectQuery,
SmartList, TraceKind, TraceNode, UpdateCommand,
};
use teaql_data_service::{MutationRequest, QueryRequest};
use crate::{DataServiceError, MetadataStore, RuntimeError};
use super::RuntimeDataService;
fn conflict_id(id: &teaql_core::Value) -> String {
id.try_u64()
.map(|value| value.to_string())
.unwrap_or_else(|| format!("{id:?}"))
}
fn sql_statement_trace(mut lineage: Vec<TraceNode>, entity: &str) -> Vec<TraceNode> {
if !lineage
.iter()
.any(|frame| frame.kind == TraceKind::Entity && frame.entity_type == entity)
{
lineage.push(TraceNode::typed(TraceKind::Entity, entity, None, ""));
}
lineage
}
impl<'a, M, E> RuntimeDataService<'a, M, E>
where
M: MetadataStore,
E: teaql_data_service::QueryExecutor + teaql_data_service::MutationExecutor + Sync,
{
pub(crate) fn new(metadata: &'a M, executor: &'a E) -> Self {
Self { metadata, executor }
}
pub(crate) async fn fetch_all(
&self,
query: &SelectQuery,
) -> Result<Vec<CompactRow>, DataServiceError<E::Error>> {
let request = QueryRequest {
query: query.clone(),
trace_chain: query.trace_chain.clone(),
comment: query.comment.clone(),
capture_debug_query: self.metadata.capture_query_debug(),
capture_execution_metadata: self.metadata.capture_execution_metadata(),
};
let res = self
.executor
.query(request)
.await
.map_err(DataServiceError::Executor)?;
Ok(res.rows)
}
pub(crate) async fn fetch_smart_list(
&self,
query: &SelectQuery,
) -> Result<SmartList<CompactRow>, DataServiceError<E::Error>> {
let request = QueryRequest {
query: query.clone(),
trace_chain: query.trace_chain.clone(),
comment: query.comment.clone(),
capture_debug_query: self.metadata.capture_query_debug(),
capture_execution_metadata: self.metadata.capture_execution_metadata(),
};
let res = self
.executor
.query(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
Ok(SmartList::from(res.rows))
}
pub(crate) async fn fetch_entities<T>(
&self,
query: &SelectQuery,
) -> Result<SmartList<T>, DataServiceError<E::Error>>
where
T: Entity,
{
let request = QueryRequest {
query: query.clone(),
trace_chain: query.trace_chain.clone(),
comment: query.comment.clone(),
capture_debug_query: self.metadata.capture_query_debug(),
capture_execution_metadata: self.metadata.capture_execution_metadata(),
};
let result = self
.executor
.query(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&result.metadata);
let entities = result
.rows
.into_iter()
.map(T::from_compact_row)
.collect::<Result<Vec<_>, _>>();
entities
.map(SmartList::from)
.map_err(DataServiceError::Entity)
}
pub(crate) async fn fetch_enhanced_entities<T>(
&self,
query: &SelectQuery,
) -> Result<SmartList<T>, DataServiceError<E::Error>>
where
T: Entity,
{
self.fetch_entities(query).await
}
pub(crate) async fn insert(
&self,
command: &InsertCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut command = command.clone();
command.trace_chain = sql_statement_trace(command.trace_chain, &command.entity);
let request = MutationRequest::Insert(command);
let res = self
.executor
.mutate(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
Ok(res.affected_rows)
}
pub(crate) async fn update(
&self,
command: &UpdateCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut sql_command = command.clone();
sql_command.trace_chain = sql_statement_trace(sql_command.trace_chain, &sql_command.entity);
let request = MutationRequest::Update(sql_command);
let res = self
.executor
.mutate(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
let affected = res.affected_rows;
if command.expected_version.is_some() && affected == 0 {
return Err(DataServiceError::Runtime(
RuntimeError::OptimisticLockConflict {
entity: command.entity.clone(),
id: conflict_id(&command.id),
},
));
}
Ok(affected)
}
pub(crate) async fn delete(
&self,
command: &DeleteCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut sql_command = command.clone();
sql_command.trace_chain = sql_statement_trace(sql_command.trace_chain, &sql_command.entity);
let request = MutationRequest::Delete(sql_command);
let res = self
.executor
.mutate(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
let affected = res.affected_rows;
if command.expected_version.is_some() && affected == 0 {
return Err(DataServiceError::Runtime(
RuntimeError::OptimisticLockConflict {
entity: command.entity.clone(),
id: conflict_id(&command.id),
},
));
}
Ok(affected)
}
pub(crate) async fn batch_insert(
&self,
command: &teaql_core::BatchInsertCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut affected = 0;
for (i, val) in command.batch_values.iter().enumerate() {
let mut insert_cmd = InsertCommand::new(command.entity.clone());
insert_cmd.values = val.clone();
if i < command.trace_chains.len() {
insert_cmd.trace_chain = command.trace_chains[i].clone();
}
insert_cmd.trace_chain =
sql_statement_trace(insert_cmd.trace_chain, &insert_cmd.entity);
let res = self
.executor
.mutate(MutationRequest::Insert(insert_cmd))
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
affected += res.affected_rows;
}
Ok(affected)
}
pub(crate) async fn batch_update(
&self,
command: &teaql_core::BatchUpdateCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut affected = 0;
for (i, val) in command.batch_values.iter().enumerate() {
let mut update_cmd =
UpdateCommand::new(command.entity.clone(), command.batch_ids[i].clone());
let mut filtered_values = MutationValues::new();
for field in &command.update_fields {
if let Some(v) = val.get(field) {
filtered_values.insert(field.clone(), v.clone());
}
}
update_cmd.values = filtered_values;
if let Some(Some(v)) = command.batch_expected_versions.get(i) {
update_cmd.expected_version = Some(*v);
}
if let Some(old) = command.batch_old_values.get(i) {
update_cmd.old_values = old.clone();
}
if i < command.trace_chains.len() {
update_cmd.trace_chain = command.trace_chains[i].clone();
}
update_cmd.trace_chain =
sql_statement_trace(update_cmd.trace_chain, &update_cmd.entity);
let expected_version = update_cmd.expected_version;
let res = self
.executor
.mutate(MutationRequest::Update(update_cmd))
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
if expected_version.is_some() && res.affected_rows == 0 {
return Err(DataServiceError::Runtime(
RuntimeError::OptimisticLockConflict {
entity: command.entity.clone(),
id: conflict_id(&command.batch_ids[i]),
},
));
}
affected += res.affected_rows;
}
if command.batch_expected_versions.iter().any(|v| v.is_some())
&& affected != command.batch_ids.len() as u64
{
return Err(DataServiceError::Runtime(RuntimeError::Graph(format!(
"batch update for {} affected {} rows for {} IDs despite version checks",
command.entity,
affected,
command.batch_ids.len()
))));
}
Ok(affected)
}
pub(crate) async fn recover(
&self,
command: &RecoverCommand,
) -> Result<u64, DataServiceError<E::Error>> {
let mut sql_command = command.clone();
sql_command.trace_chain = sql_statement_trace(sql_command.trace_chain, &sql_command.entity);
let request = MutationRequest::Recover(sql_command);
let res = self
.executor
.mutate(request)
.await
.map_err(DataServiceError::Executor)?;
self.metadata.record_metadata_log(&res.metadata);
let affected = res.affected_rows;
if affected == 0 {
return Err(DataServiceError::Runtime(
RuntimeError::OptimisticLockConflict {
entity: command.entity.clone(),
id: format!("{:?}", command.id),
},
));
}
Ok(affected)
}
pub(crate) async fn insert_many(
&self,
commands: &[InsertCommand],
) -> Result<u64, DataServiceError<E::Error>> {
let mut total = 0;
for command in commands {
total += self.insert(command).await?;
}
Ok(total)
}
pub(crate) async fn update_many(
&self,
commands: &[UpdateCommand],
) -> Result<u64, DataServiceError<E::Error>> {
let mut total = 0;
for command in commands {
total += self.update(command).await?;
}
Ok(total)
}
pub(crate) async fn delete_many(
&self,
commands: &[DeleteCommand],
) -> Result<u64, DataServiceError<E::Error>> {
let mut total = 0;
for command in commands {
total += self.delete(command).await?;
}
Ok(total)
}
pub(crate) async fn recover_many(
&self,
commands: &[RecoverCommand],
) -> Result<u64, DataServiceError<E::Error>> {
let mut total = 0;
for command in commands {
total += self.recover(command).await?;
}
Ok(total)
}
}