use std::any::Any;
use std::collections::BTreeMap;
use std::sync::{Arc, Mutex};
use std::time::Instant;
use anyhow::Result;
use lora_compiler::CompiledQuery;
use lora_executor::{LoraValue, MutableExecutionContext, MutableExecutor, Row};
use lora_store::{GraphStorage, GraphStorageMut, MutationEvent, MutationRecorder};
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;
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>> {
let _commit_lock = self
.writer
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
let buffer = Arc::new(Mutex::new(Vec::<MutationEvent>::new()));
let buffering_rec: Arc<dyn MutationRecorder> =
Arc::new(BufferingRecorder::new(buffer.clone()));
let mut handle = self.store.write();
let exec_result = {
let staged = handle.as_mut();
install_recorder_if_inmemory(staged, Some(buffering_rec));
let mut executor = MutableExecutor::with_deadline(
MutableExecutionContext {
storage: staged,
params,
},
deadline,
);
let r = executor.execute_compiled_rows(compiled);
install_recorder_if_inmemory(staged, None);
r
};
let rows = match exec_result {
Ok(rows) => rows,
Err(e) => return Err(anyhow::Error::from(e)),
};
let events: Vec<MutationEvent> = std::mem::take(&mut buffer.lock().unwrap());
if events.is_empty() {
if let Some(rec) = self.wal.as_ref() {
let staged = handle.as_mut();
install_recorder_if_inmemory(
staged,
Some(rec.clone() as Arc<dyn MutationRecorder>),
);
}
return Ok(rows);
}
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() {
let staged = handle.as_mut();
install_recorder_if_inmemory(staged, Some(rec.clone() as Arc<dyn MutationRecorder>));
}
if wrote_commit {
if let Some(rec) = self.wal.as_ref() {
let live = handle.snapshot();
self.observe_snapshot_commit_if_needed(&*live, rec)?;
}
}
Ok(rows)
}
}