use std::any::Any;
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use anyhow::{anyhow, Result};
use lora_compiler::CompiledQuery;
use lora_executor::{LoraValue, MutableExecutionContext, MutableExecutor, Row};
use lora_store::{
GraphStorage, GraphStorageMut, MutationEvent, MutationRecorder, MutationWriteSet,
};
use crate::database::Database;
use crate::transaction::BufferingRecorder;
use crate::wal::write_scope::ensure_wal_query_can_start;
use super::replay::{
install_recorder_if_inmemory, merge_events_into, validate_write_set_unchanged,
};
impl<S> Database<S>
where
S: GraphStorage + GraphStorageMut + Any + Clone + Send + Sync + 'static,
{
pub(crate) fn execute_mutating_optimistic(
&self,
params: BTreeMap<String, LoraValue>,
deadline: Option<Instant>,
compiled: &Arc<CompiledQuery>,
) -> Result<Vec<Row>> {
const MAX_RETRIES: usize = 64;
for _ in 0..MAX_RETRIES {
let snapshot = self.store.load_full();
let mut staged: S = (*snapshot).clone();
let buffer = Arc::new(Mutex::new(Vec::<MutationEvent>::new()));
let buffering_rec: Arc<dyn MutationRecorder> =
Arc::new(BufferingRecorder::new(buffer.clone()));
install_recorder_if_inmemory(&mut staged, Some(buffering_rec));
let exec_result = {
let mut executor = MutableExecutor::with_deadline(
MutableExecutionContext {
storage: &mut staged,
params: params.clone(),
},
deadline,
);
executor.execute_compiled_rows(compiled)
};
let rows = match exec_result {
Ok(rows) => rows,
Err(e) => return Err(anyhow::Error::from(e)),
};
install_recorder_if_inmemory(&mut staged, None);
let events: Vec<MutationEvent> = std::mem::take(&mut buffer.lock().unwrap());
if events.is_empty() {
return Ok(rows);
}
let mut write_set = MutationWriteSet::new();
write_set.extend_from_events(events.iter());
let _commit_lock = self
.writer
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let current = self.store.load_full();
let publish_state: S = if Arc::ptr_eq(¤t, &snapshot) {
staged
} else {
if !validate_write_set_unchanged(&*snapshot, &*current, &write_set) {
drop(_commit_lock);
continue; }
let mut merged: S = (*current).clone();
if !merge_events_into(&mut merged, &events) {
drop(_commit_lock);
continue;
}
merged
};
let mut publish_state = publish_state;
let mut wrote_commit = false;
if let Some(rec) = self.wal.as_ref() {
ensure_wal_query_can_start(rec)?;
wrote_commit = rec.commit_events(events)?.wrote();
}
if let Some(rec) = self.wal.as_ref() {
install_recorder_if_inmemory(
&mut publish_state,
Some(rec.clone() as Arc<dyn MutationRecorder>),
);
}
self.store.store(Arc::new(publish_state));
if wrote_commit {
if let Some(rec) = self.wal.as_ref() {
let live = self.store.load_full();
self.observe_snapshot_commit_if_needed(&*live, rec)?;
}
}
return Ok(rows);
}
Err(anyhow!(
"auto-commit write conflict: exceeded {MAX_RETRIES} retries"
))
}
}