1use std::collections::BTreeMap;
2use std::sync::atomic::{AtomicBool, Ordering};
3use std::sync::{Arc, Mutex, MutexGuard};
4use std::time::Duration;
5
6use web_time::Instant;
7
8use anyhow::Result;
9use thiserror::Error;
10
11use lora_analyzer::Analyzer;
12use lora_compiler::{CompiledQuery, Compiler};
13use lora_executor::{
14 classify_stream, compiled_result_columns, project_rows, ExecuteOptions, ExecutionContext,
15 Executor, LoraValue, MutableExecutionContext, MutableExecutor, MutablePullExecutor,
16 PullExecutor, QueryResult, Row, RowSource,
17};
18use lora_parser::parse_query;
19use lora_store::{InMemoryGraph, MutationEvent, MutationRecorder};
20use lora_wal::WalRecorder;
21
22use crate::changes::ChangeHub;
23use crate::error::LoraError;
24use crate::explain::{OperatorMetrics, ProfileMetrics, QueryPlan, QueryProfile};
25use crate::live_store::LiveStore;
26use crate::snapshot::ManagedSnapshotStore;
27use crate::stream::QueryStream;
28use crate::wal::write_scope::ensure_wal_not_poisoned;
29use lora_compiler::plan_tree_from_compiled;
30use lora_executor::{plan_result_columns, CollectorGuard, MetricsCollector};
31
32#[derive(Debug, Clone, PartialEq, Eq, Error)]
38pub enum TransactionError {
39 #[error("transaction is already closed")]
40 AlreadyClosed,
41
42 #[error("transaction has no live graph guard")]
43 NoGraphGuard,
44
45 #[error("transaction has no staged graph")]
46 NoStagedGraph,
47
48 #[error("cannot commit transaction while a streaming cursor is still active")]
49 CursorActiveCommit,
50
51 #[error("cannot start a new statement while a streaming cursor is still active")]
52 CursorActiveStatement,
53
54 #[error("cannot execute mutating query in read-only transaction")]
55 ReadOnlyMutation,
56
57 #[error("streaming write cursor requires a ReadWrite transaction")]
58 StreamingRequiresReadWrite,
59
60 #[error("read-only transaction cannot publish staged graph")]
61 ReadOnlyCommit,
62
63 #[error("transaction state lock is poisoned")]
64 Poisoned,
65}
66
67#[derive(Debug, Clone, Copy, PartialEq, Eq)]
69pub enum TransactionMode {
70 ReadOnly,
72 ReadWrite,
74}
75
76pub(crate) enum LiveStoreGuard<'db> {
86 Read(Arc<InMemoryGraph>),
87 Write(WriteLease<'db>),
88}
89
90pub(crate) struct WriteLease<'db> {
95 pub(crate) _writer_lock: MutexGuard<'db, ()>,
97 pub(crate) store: Arc<LiveStore<InMemoryGraph>>,
99 pub(crate) snapshot: Arc<InMemoryGraph>,
102}
103
104impl LiveStoreGuard<'_> {
105 fn as_graph(&self) -> &InMemoryGraph {
106 match self {
107 Self::Read(arc) => arc,
108 Self::Write(lease) => &lease.snapshot,
109 }
110 }
111}
112
113pub(crate) struct Savepoint {
118 staged: Option<InMemoryGraph>,
119 buffer_len: usize,
120}
121
122#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126pub(crate) enum TxStreamOutcome {
127 Exhausted,
128 Interrupted,
129}
130
131impl TxStreamOutcome {
132 fn should_restore_savepoint(self, rollback_on_drop: bool) -> bool {
133 matches!(self, Self::Interrupted) && rollback_on_drop
134 }
135}
136
137pub(crate) struct TxCursorLease {
144 handle: Arc<Mutex<TxInner>>,
145 rollback_on_drop: bool,
146 finalized: bool,
147}
148
149impl TxCursorLease {
150 pub(crate) fn new(handle: Arc<Mutex<TxInner>>, rollback_on_drop: bool) -> Self {
151 Self {
152 handle,
153 rollback_on_drop,
154 finalized: false,
155 }
156 }
157
158 pub(crate) fn finalize(&mut self, outcome: TxStreamOutcome) {
159 if self.finalized {
160 return;
161 }
162 finalize_tx_stream(&self.handle, outcome, self.rollback_on_drop);
163 self.finalized = true;
164 }
165}
166
167impl Drop for TxCursorLease {
168 fn drop(&mut self) {
169 self.finalize(TxStreamOutcome::Interrupted);
170 }
171}
172
173pub(crate) struct BufferingRecorder {
180 buffer: Arc<Mutex<Vec<MutationEvent>>>,
181 failed: Arc<AtomicBool>,
182}
183
184impl BufferingRecorder {
185 pub(crate) fn new(buffer: Arc<Mutex<Vec<MutationEvent>>>, failed: Arc<AtomicBool>) -> Self {
186 Self { buffer, failed }
187 }
188}
189
190impl MutationRecorder for BufferingRecorder {
191 fn record(&self, event: MutationEvent) {
192 match self.buffer.lock() {
193 Ok(mut buf) => buf.push(event),
194 Err(_) => self.failed.store(true, Ordering::Release),
195 }
196 }
197}
198
199pub(crate) struct TxInner {
204 pub(crate) staged: Option<InMemoryGraph>,
208 pub(crate) buffer: Arc<Mutex<Vec<MutationEvent>>>,
212 pub(crate) buffer_failed: Arc<AtomicBool>,
216 pub(crate) pending_savepoint: Option<Savepoint>,
220 pub(crate) cursor_active: bool,
224 pub(crate) closed: bool,
228 pub(crate) mode: TransactionMode,
230 pub(crate) buffer_mutations: bool,
234}
235
236pub struct Transaction<'db> {
252 pub(crate) live: Option<LiveStoreGuard<'db>>,
253 pub(crate) inner: Arc<Mutex<TxInner>>,
254 pub(crate) wal: Option<Arc<WalRecorder>>,
255 pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
256 pub(crate) changes: Arc<ChangeHub>,
257 mode: TransactionMode,
258}
259
260impl<'db> Transaction<'db> {
261 pub(crate) fn new(
263 live: LiveStoreGuard<'db>,
264 wal: Option<Arc<WalRecorder>>,
265 snapshots: Option<Arc<ManagedSnapshotStore>>,
266 changes: Arc<ChangeHub>,
267 mode: TransactionMode,
268 ) -> Self {
269 let buffer_mutations = wal.is_some() || changes.is_active();
272 let inner = TxInner {
273 staged: None,
274 buffer: Arc::new(Mutex::new(Vec::new())),
275 buffer_failed: Arc::new(AtomicBool::new(false)),
276 pending_savepoint: None,
277 cursor_active: false,
278 closed: false,
279 mode,
280 buffer_mutations,
281 };
282 Self {
283 live: Some(live),
284 inner: Arc::new(Mutex::new(inner)),
285 wal,
286 snapshots,
287 changes,
288 mode,
289 }
290 }
291
292 pub fn mode(&self) -> TransactionMode {
294 self.mode
295 }
296
297 pub fn execute(
300 &mut self,
301 query: &str,
302 options: Option<ExecuteOptions>,
303 ) -> Result<QueryResult, LoraError> {
304 self.execute_with_params(query, options, BTreeMap::new())
305 }
306
307 pub fn execute_with_timeout(
309 &mut self,
310 query: &str,
311 options: Option<ExecuteOptions>,
312 timeout: Duration,
313 ) -> Result<QueryResult, LoraError> {
314 let rows = self.execute_rows_with_params_deadline(
315 query,
316 BTreeMap::new(),
317 Some(deadline_after(timeout)),
318 )?;
319 Ok(project_rows(rows, options.unwrap_or_default()))
320 }
321
322 pub fn execute_with_params(
324 &mut self,
325 query: &str,
326 options: Option<ExecuteOptions>,
327 params: BTreeMap<String, LoraValue>,
328 ) -> Result<QueryResult, LoraError> {
329 let rows = self.execute_rows_with_params_deadline(query, params, None)?;
330 Ok(project_rows(rows, options.unwrap_or_default()))
331 }
332
333 pub fn execute_with_params_timeout(
336 &mut self,
337 query: &str,
338 options: Option<ExecuteOptions>,
339 params: BTreeMap<String, LoraValue>,
340 timeout: Duration,
341 ) -> Result<QueryResult, LoraError> {
342 let rows =
343 self.execute_rows_with_params_deadline(query, params, Some(deadline_after(timeout)))?;
344 Ok(project_rows(rows, options.unwrap_or_default()))
345 }
346
347 pub fn execute_with_params_deadline(
352 &mut self,
353 query: &str,
354 options: Option<ExecuteOptions>,
355 params: BTreeMap<String, LoraValue>,
356 deadline: Instant,
357 ) -> Result<QueryResult, LoraError> {
358 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
359 Ok(project_rows(rows, options.unwrap_or_default()))
360 }
361
362 pub fn execute_rows(&mut self, query: &str) -> Result<Vec<Row>, LoraError> {
365 self.execute_rows_with_params(query, BTreeMap::new())
366 }
367
368 pub fn execute_rows_with_params(
371 &mut self,
372 query: &str,
373 params: BTreeMap<String, LoraValue>,
374 ) -> Result<Vec<Row>, LoraError> {
375 Ok(self.execute_rows_with_params_deadline(query, params, None)?)
376 }
377
378 fn execute_rows_with_params_deadline(
379 &mut self,
380 query: &str,
381 params: BTreeMap<String, LoraValue>,
382 deadline: Option<Instant>,
383 ) -> Result<Vec<Row>> {
384 if crate::Database::<InMemoryGraph>::is_schema_command_text(query) {
385 let document = parse_query(query)?;
386 if let lora_ast::Statement::Schema(command) = &document.statement {
387 return self.execute_schema_in_tx(command, ¶ms);
388 }
389 }
390 let compiled = self.compile_in_tx(query)?;
391 self.execute_rows_compiled_deadline(&compiled, params, deadline)
392 }
393
394 fn execute_rows_compiled_deadline(
395 &mut self,
396 compiled: &CompiledQuery,
397 params: BTreeMap<String, LoraValue>,
398 deadline: Option<Instant>,
399 ) -> Result<Vec<Row>> {
400 if self.is_read_only_unchecked() {
402 self.precheck_open_no_savepoint()?;
403 return self.execute_live_compiled(compiled, params, deadline);
404 }
405
406 let mut inner = self.begin_statement()?;
408 let is_mutating = classify_stream(compiled).is_mutating();
409
410 if !is_mutating {
411 return self.execute_read_statement(inner, compiled, params, deadline);
417 }
418
419 let savepoint = self.prepare_mutating_statement(&mut inner)?;
423
424 let exec_result: ExecResultRows = {
425 let staged = inner.staged_mut()?;
426 execute_mutable_compiled(staged, compiled, params, deadline)
427 };
428
429 match exec_result {
430 Ok(rows) => Ok(rows),
431 Err(err) => {
432 restore_savepoint(&mut inner, savepoint);
433 Err(err)
434 }
435 }
436 }
437
438 pub(crate) fn open_streaming_compiled_autocommit(
467 &mut self,
468 compiled: Arc<CompiledQuery>,
469 params: BTreeMap<String, LoraValue>,
470 ) -> Result<Box<dyn RowSource + 'static>> {
471 if self.is_read_only_unchecked() {
472 return Err(TransactionError::StreamingRequiresReadWrite.into());
473 }
474
475 let mut inner = self.begin_statement()?;
476 self.ensure_staged_locked(&mut inner)?;
477 inner.activate_cursor();
478
479 let staged_ptr: *mut InMemoryGraph = inner
483 .staged
484 .as_mut()
485 .ok_or(TransactionError::NoStagedGraph)?
486 as *mut _;
487 drop(inner);
488
489 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
495 let compiled_static: &'static CompiledQuery =
496 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
497
498 let cursor = MutablePullExecutor::new(storage_static, params)
502 .open_compiled(compiled_static)
503 .map_err(|e| {
504 if let Ok(mut inner) = self.inner.lock() {
508 discard_transaction_state(&mut inner);
509 }
510 self.live.take();
511 anyhow::Error::from(e)
512 })?;
513
514 Ok(Box::new(StreamingCursorWithArc {
520 cursor,
521 _compiled: compiled,
522 }))
523 }
524
525 fn compile_in_tx(&self, query: &str) -> Result<CompiledQuery> {
531 let document = parse_query(query)?;
532 let (resolved, stats) = {
533 let inner = self.lock_inner()?;
534 if let Some(staged) = &inner.staged {
535 let mut analyzer = Analyzer::new(staged);
536 let resolved = analyzer.analyze(&document)?;
537 let stats = staged.graph_stats();
538 (resolved, stats)
539 } else {
540 drop(inner);
541 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
542 let graph = live.as_graph();
543 let mut analyzer = Analyzer::new(graph);
544 let resolved = analyzer.analyze(&document)?;
545 let stats = graph.graph_stats();
546 (resolved, stats)
547 }
548 };
549 Ok(Compiler::compile(&resolved, &stats))
550 }
551
552 fn ensure_staged_locked(&self, inner: &mut MutexGuard<'_, TxInner>) -> Result<()> {
556 if inner.staged.is_some() {
557 return Ok(());
558 }
559 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
560 let mut staged: InMemoryGraph = live.as_graph().clone();
561 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
562 staged.set_mutation_recorder(Some(Arc::new(BufferingRecorder::new(
563 inner.buffer.clone(),
564 inner.buffer_failed.clone(),
565 )) as Arc<dyn MutationRecorder>));
566 }
567 inner.staged = Some(staged);
568 Ok(())
569 }
570
571 fn execute_live_compiled(
572 &self,
573 compiled: &CompiledQuery,
574 params: BTreeMap<String, LoraValue>,
575 deadline: Option<Instant>,
576 ) -> Result<Vec<Row>> {
577 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
578 execute_read_compiled(live.as_graph(), compiled, params, deadline)
579 }
580
581 fn execute_read_statement(
582 &self,
583 inner: MutexGuard<'_, TxInner>,
584 compiled: &CompiledQuery,
585 params: BTreeMap<String, LoraValue>,
586 deadline: Option<Instant>,
587 ) -> Result<Vec<Row>> {
588 match inner.staged.as_ref() {
589 Some(staged) => execute_read_compiled(staged, compiled, params, deadline),
590 None => {
591 drop(inner);
592 self.execute_live_compiled(compiled, params, deadline)
593 }
594 }
595 }
596
597 fn execute_schema_in_tx(
604 &mut self,
605 command: &lora_ast::SchemaCommand,
606 params: &BTreeMap<String, LoraValue>,
607 ) -> Result<Vec<Row>> {
608 use crate::database::schema::{apply_schema_mutation, schema_command_is_read, show_schema};
609
610 if schema_command_is_read(command) {
611 if self.is_read_only_unchecked() {
612 self.precheck_open_no_savepoint()?;
613 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
614 return show_schema(live.as_graph(), command, params);
615 }
616 let inner = self.begin_statement()?;
617 if let Some(staged) = inner.staged.as_ref() {
618 return show_schema(staged, command, params);
619 }
620 drop(inner);
621 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
622 return show_schema(live.as_graph(), command, params);
623 }
624
625 if self.is_read_only_unchecked() {
626 return Err(TransactionError::ReadOnlyMutation.into());
627 }
628 let mut inner = self.begin_statement()?;
629 let savepoint = self.prepare_mutating_statement(&mut inner)?;
630 let result = {
631 let staged = inner.staged_mut()?;
632 apply_schema_mutation(staged, command, params)
633 };
634 if result.is_err() {
635 restore_savepoint(&mut inner, savepoint);
636 }
637 result
638 }
639
640 fn prepare_mutating_statement(
641 &self,
642 inner: &mut MutexGuard<'_, TxInner>,
643 ) -> Result<Option<Savepoint>> {
644 let clone_savepoint_graph = inner.staged.is_some();
645 self.ensure_staged_locked(inner)?;
646 Ok(Some(take_savepoint(inner, clone_savepoint_graph)))
647 }
648
649 pub fn explain(
654 &self,
655 query: &str,
656 _params: Option<BTreeMap<String, LoraValue>>,
657 ) -> Result<QueryPlan, LoraError> {
658 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
659 Ok(query_plan_for(query, &compiled))
660 }
661
662 pub fn profile(
670 &mut self,
671 query: &str,
672 params: Option<BTreeMap<String, LoraValue>>,
673 ) -> Result<QueryProfile, LoraError> {
674 let params = params.unwrap_or_default();
675 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
676 let plan = query_plan_for(query, &compiled);
677 let shape = plan.shape;
678
679 let collector = Arc::new(MetricsCollector::new());
680 let _guard = CollectorGuard::install(collector.clone());
681
682 let started = Instant::now();
683 let rows = self
684 .execute_rows_compiled_deadline(&compiled, params, None)
685 .map_err(LoraError::from_anyhow)?;
686 let total_elapsed_ns = started.elapsed().as_nanos() as u64;
687
688 drop(_guard);
689 let per_operator = collector
690 .snapshot()
691 .into_iter()
692 .map(|(id, op)| {
693 (
694 id,
695 OperatorMetrics {
696 rows: op.rows,
697 elapsed_ns: op.elapsed_ns,
698 next_calls: op.next_calls,
699 db_hits: 0,
700 },
701 )
702 })
703 .collect();
704
705 let metrics = ProfileMetrics {
706 total_elapsed_ns,
707 total_rows: rows.len() as u64,
708 mutated: shape.is_mutating(),
709 per_operator,
710 };
711
712 Ok(QueryProfile { plan, metrics })
713 }
714
715 pub fn stream(&mut self, query: &str) -> Result<QueryStream<'static>, LoraError> {
717 self.stream_with_params(query, BTreeMap::new())
718 }
719
720 pub fn stream_with_params(
723 &mut self,
724 query: &str,
725 params: BTreeMap<String, LoraValue>,
726 ) -> Result<QueryStream<'static>, LoraError> {
727 let compiled = Arc::new(self.compile_in_tx(query)?);
728 let columns = compiled_result_columns(&compiled);
729 Ok(self.stream_compiled(compiled, columns, params)?)
730 }
731
732 pub(crate) fn stream_compiled(
736 &mut self,
737 compiled: Arc<CompiledQuery>,
738 columns: Vec<String>,
739 params: BTreeMap<String, LoraValue>,
740 ) -> Result<QueryStream<'static>> {
741 let mut inner = self.begin_statement()?;
742 let is_mutating = classify_stream(&compiled).is_mutating();
743 ensure_stream_allowed(&inner, is_mutating)?;
744
745 let rollback_on_drop = stream_rolls_back_on_drop(is_mutating);
746 let staged_ptr = self.prepare_stream_staging(&mut inner, rollback_on_drop)?;
747 drop(inner);
748
749 let cursor = open_tx_stream_cursor(staged_ptr, compiled, params, is_mutating);
750
751 match cursor {
752 Ok(cursor) => Ok(QueryStream::for_tx_cursor(
753 cursor,
754 columns,
755 TxCursorLease::new(self.inner.clone(), rollback_on_drop),
756 )),
757 Err(err) => {
758 finalize_tx_stream(&self.inner, TxStreamOutcome::Interrupted, rollback_on_drop);
759 Err(err)
760 }
761 }
762 }
763
764 fn prepare_stream_staging(
765 &self,
766 inner: &mut MutexGuard<'_, TxInner>,
767 rollback_on_drop: bool,
768 ) -> Result<*mut InMemoryGraph> {
769 let clone_savepoint_graph = inner.staged.is_some();
774 self.ensure_staged_locked(inner)?;
775 inner.activate_cursor();
776
777 if rollback_on_drop {
778 inner.pending_savepoint = Some(take_savepoint(inner, clone_savepoint_graph));
779 } else {
780 inner.pending_savepoint = None;
781 }
782
783 Ok(inner
784 .staged
785 .as_mut()
786 .ok_or(TransactionError::NoStagedGraph)? as *mut _)
787 }
788
789 pub fn commit(mut self) -> Result<(), LoraError> {
796 let CommitState {
797 staged,
798 buffer_events,
799 mode,
800 } = self.take_commit_state()?;
801
802 let capture = self.changes.is_active()
803 && matches!(mode, TransactionMode::ReadWrite)
804 && staged.is_some()
805 && !buffer_events.is_empty();
806 let (wrote_wal_commit, captured) = if capture && self.wal.is_some() {
807 let events = buffer_events.clone();
808 let lsn = self.replay_commit_wal_lsn(mode, buffer_events)?;
809 (lsn.is_some(), lsn.map(|lsn| (Some(lsn), events)))
810 } else if capture {
811 (false, Some((None, buffer_events)))
812 } else {
813 (self.replay_commit_wal(mode, buffer_events)?, None)
814 };
815 self.publish_staged_graph(mode, staged, wrote_wal_commit)?;
816
817 if let Some((lsn, events)) = captured {
818 if let Some(LiveStoreGuard::Write(lease)) = &self.live {
821 let post = lease.store.load_full();
822 crate::changes::publish_committed(
823 &self.changes,
824 lsn.map(|lsn| lsn.raw()),
825 &events,
826 &crate::changes::PreImages::for_events(&events, Some(&lease.snapshot)),
827 &post,
828 );
829 }
830 }
831
832 self.live.take();
833 Ok(())
834 }
835
836 fn take_commit_state(&self) -> Result<CommitState> {
837 let mut inner = self.lock_inner()?;
838 if inner.cursor_active {
839 return Err(TransactionError::CursorActiveCommit.into());
840 }
841 if inner.closed {
842 return Err(TransactionError::AlreadyClosed.into());
843 }
844
845 let mode = inner.mode;
846 if inner.buffer_failed.load(Ordering::Acquire) {
847 inner.closed = true;
848 return Err(TransactionError::Poisoned.into());
849 }
850 let buffer_events = {
851 let mut buffer = inner
852 .buffer
853 .lock()
854 .map_err(|_| TransactionError::Poisoned)?;
855 std::mem::take(&mut *buffer)
856 };
857 let staged = inner.staged.take();
861 inner.closed = true;
862
863 Ok(CommitState {
864 staged,
865 buffer_events,
866 mode,
867 })
868 }
869
870 fn replay_commit_wal(
871 &self,
872 mode: TransactionMode,
873 buffer_events: Vec<MutationEvent>,
874 ) -> Result<bool> {
875 let Some(rec) = &self.wal else {
876 return Ok(false);
877 };
878
879 if !matches!(mode, TransactionMode::ReadWrite) {
880 ensure_wal_not_poisoned(rec)?;
881 return Ok(false);
882 }
883
884 Ok(rec.commit_events(buffer_events)?.wrote())
885 }
886
887 fn replay_commit_wal_lsn(
890 &self,
891 mode: TransactionMode,
892 buffer_events: Vec<MutationEvent>,
893 ) -> Result<Option<lora_wal::Lsn>> {
894 let Some(rec) = &self.wal else {
895 return Ok(None);
896 };
897 if !matches!(mode, TransactionMode::ReadWrite) {
898 ensure_wal_not_poisoned(rec)?;
899 return Ok(None);
900 }
901 Ok(rec.commit_events_lsn(buffer_events)?)
902 }
903
904 fn publish_staged_graph(
905 &mut self,
906 mode: TransactionMode,
907 staged: Option<InMemoryGraph>,
908 wrote_wal_commit: bool,
909 ) -> Result<()> {
910 if !matches!(mode, TransactionMode::ReadWrite) {
911 return Ok(());
912 }
913
914 let Some(mut staged) = staged else {
915 return Ok(());
916 };
917
918 staged.set_mutation_recorder(None);
923 let wal = self.wal.clone();
924 if let Some(rec) = &wal {
925 staged.set_mutation_recorder(Some(rec.clone() as Arc<dyn MutationRecorder>));
926 }
927
928 let live = self.live.as_mut().ok_or(TransactionError::NoGraphGuard)?;
929 let lease = match live {
930 LiveStoreGuard::Write(lease) => lease,
931 LiveStoreGuard::Read(_) => {
932 return Err(TransactionError::ReadOnlyCommit.into());
933 }
934 };
935
936 if wrote_wal_commit {
937 if let (Some(snapshots), Some(rec)) = (&self.snapshots, wal.as_ref()) {
938 snapshots.observe_commit(&staged, rec)?;
939 }
940 }
941
942 lease.store.store(Arc::new(staged));
946
947 Ok(())
948 }
949
950 pub fn rollback(mut self) -> Result<(), LoraError> {
953 let mut inner = self.lock_inner()?;
954 if inner.closed {
955 return Err(TransactionError::AlreadyClosed.into());
956 }
957 discard_transaction_state(&mut inner);
958 drop(inner);
959 self.live.take();
960 Ok(())
961 }
962
963 fn begin_statement(&self) -> Result<MutexGuard<'_, TxInner>> {
969 let inner = self.lock_inner()?;
970 if inner.closed {
971 return Err(TransactionError::AlreadyClosed.into());
972 }
973 if inner.cursor_active {
974 return Err(TransactionError::CursorActiveStatement.into());
975 }
976 Ok(inner)
977 }
978
979 fn precheck_open_no_savepoint(&self) -> Result<()> {
983 let inner = self.lock_inner()?;
984 if inner.closed {
985 return Err(TransactionError::AlreadyClosed.into());
986 }
987 if inner.cursor_active {
988 return Err(TransactionError::CursorActiveStatement.into());
989 }
990 Ok(())
991 }
992
993 fn is_read_only_unchecked(&self) -> bool {
997 matches!(self.mode, TransactionMode::ReadOnly)
998 }
999
1000 fn lock_inner(&self) -> Result<MutexGuard<'_, TxInner>> {
1001 self.inner
1002 .lock()
1003 .map_err(|_| TransactionError::Poisoned.into())
1004 }
1005
1006 pub(crate) fn release_streaming_cursor(&self) {
1007 if let Ok(mut inner) = self.inner.lock() {
1008 inner.release_cursor();
1009 }
1010 }
1011}
1012
1013type ExecResultRows = Result<Vec<Row>>;
1014
1015fn deadline_after(timeout: Duration) -> Instant {
1016 Instant::now()
1017 .checked_add(timeout)
1018 .unwrap_or_else(Instant::now)
1019}
1020
1021fn execute_read_compiled(
1022 storage: &InMemoryGraph,
1023 compiled: &CompiledQuery,
1024 params: BTreeMap<String, LoraValue>,
1025 deadline: Option<Instant>,
1026) -> Result<Vec<Row>> {
1027 if crate::database::pull_mode::should_collect_read_via_pull(compiled) {
1030 return lora_executor::collect_compiled_with_deadline(storage, params, compiled, deadline)
1031 .map_err(anyhow::Error::from);
1032 }
1033 let executor = Executor::with_deadline(ExecutionContext { storage, params }, deadline);
1034 executor
1035 .execute_compiled_rows(compiled)
1036 .map_err(anyhow::Error::from)
1037}
1038
1039fn execute_mutable_compiled(
1040 storage: &mut InMemoryGraph,
1041 compiled: &CompiledQuery,
1042 params: BTreeMap<String, LoraValue>,
1043 deadline: Option<Instant>,
1044) -> Result<Vec<Row>> {
1045 let mut executor =
1046 MutableExecutor::with_deadline(MutableExecutionContext { storage, params }, deadline);
1047 executor
1048 .execute_compiled_rows(compiled)
1049 .map_err(anyhow::Error::from)
1050}
1051
1052fn query_plan_for(query: &str, compiled: &CompiledQuery) -> QueryPlan {
1053 QueryPlan {
1054 query: query.to_string(),
1055 tree: plan_tree_from_compiled(compiled),
1056 shape: classify_stream(compiled).into(),
1057 result_columns: plan_result_columns(&compiled.physical),
1058 }
1059}
1060
1061fn ensure_stream_allowed(inner: &TxInner, is_mutating: bool) -> Result<()> {
1062 if matches!(inner.mode, TransactionMode::ReadOnly) && is_mutating {
1063 Err(TransactionError::ReadOnlyMutation.into())
1064 } else {
1065 Ok(())
1066 }
1067}
1068
1069fn stream_rolls_back_on_drop(is_mutating: bool) -> bool {
1070 is_mutating
1071}
1072
1073fn open_tx_stream_cursor(
1074 staged_ptr: *mut InMemoryGraph,
1075 compiled: Arc<CompiledQuery>,
1076 params: BTreeMap<String, LoraValue>,
1077 is_mutating: bool,
1078) -> Result<Box<dyn RowSource + 'static>> {
1079 let compiled_static: &'static CompiledQuery =
1083 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
1084
1085 if is_mutating {
1086 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
1089 MutablePullExecutor::new(storage_static, params)
1090 .open_compiled(compiled_static)
1091 .map(|cursor| boxed_streaming_cursor(cursor, compiled))
1092 .map_err(anyhow::Error::from)
1093 } else {
1094 let storage_static: &'static InMemoryGraph = unsafe { &*staged_ptr };
1097 PullExecutor::new(storage_static, params)
1098 .open_compiled(compiled_static)
1099 .map(|cursor| boxed_streaming_cursor(cursor, compiled))
1100 .map_err(anyhow::Error::from)
1101 }
1102}
1103
1104fn boxed_streaming_cursor(
1105 cursor: Box<dyn RowSource + 'static>,
1106 compiled: Arc<CompiledQuery>,
1107) -> Box<dyn RowSource + 'static> {
1108 Box::new(StreamingCursorWithArc {
1109 cursor,
1110 _compiled: compiled,
1111 })
1112}
1113
1114struct CommitState {
1115 staged: Option<InMemoryGraph>,
1116 buffer_events: Vec<MutationEvent>,
1117 mode: TransactionMode,
1118}
1119
1120impl TxInner {
1121 fn staged_mut(&mut self) -> Result<&mut InMemoryGraph> {
1122 self.staged
1123 .as_mut()
1124 .ok_or(TransactionError::NoStagedGraph.into())
1125 }
1126
1127 fn activate_cursor(&mut self) {
1128 self.cursor_active = true;
1129 }
1130
1131 fn release_cursor(&mut self) {
1132 self.cursor_active = false;
1133 }
1134
1135 fn clear_pending_savepoint(&mut self) {
1136 self.pending_savepoint = None;
1137 }
1138
1139 fn restore_pending_savepoint(&mut self) {
1140 if let Some(sp) = self.pending_savepoint.take() {
1141 apply_savepoint(self, sp);
1142 }
1143 }
1144
1145 fn finalize_stream(&mut self, outcome: TxStreamOutcome, rollback_on_drop: bool) {
1146 self.release_cursor();
1147
1148 if self.closed {
1149 discard_transaction_state(self);
1150 return;
1151 }
1152
1153 if outcome.should_restore_savepoint(rollback_on_drop) {
1154 self.restore_pending_savepoint();
1155 } else {
1156 self.clear_pending_savepoint();
1157 }
1158 }
1159}
1160
1161struct StreamingCursorWithArc {
1166 cursor: Box<dyn RowSource + 'static>,
1167 _compiled: Arc<CompiledQuery>,
1168}
1169
1170impl RowSource for StreamingCursorWithArc {
1171 fn next_row(&mut self) -> lora_executor::ExecResult<Option<Row>> {
1172 self.cursor.next_row()
1173 }
1174}
1175
1176fn finalize_tx_stream(
1177 handle: &Arc<Mutex<TxInner>>,
1178 outcome: TxStreamOutcome,
1179 rollback_on_drop: bool,
1180) {
1181 if let Ok(mut inner) = handle.lock() {
1182 inner.finalize_stream(outcome, rollback_on_drop);
1183 }
1184}
1185
1186fn discard_transaction_state(inner: &mut TxInner) {
1187 inner.clear_pending_savepoint();
1189 inner.release_cursor();
1190 inner.staged = None;
1191 if let Ok(mut buf) = inner.buffer.lock() {
1192 buf.clear();
1193 } else {
1194 inner.buffer_failed.store(true, Ordering::Release);
1195 }
1196 inner.closed = true;
1197}
1198
1199fn take_savepoint(inner: &TxInner, clone_staged: bool) -> Savepoint {
1200 let buffer_len = match inner.buffer.lock() {
1201 Ok(buffer) => buffer.len(),
1202 Err(_) => {
1203 inner.buffer_failed.store(true, Ordering::Release);
1204 0
1205 }
1206 };
1207 Savepoint {
1208 staged: if clone_staged {
1209 inner.staged.as_ref().cloned()
1210 } else {
1211 None
1212 },
1213 buffer_len,
1214 }
1215}
1216
1217fn restore_savepoint(inner: &mut TxInner, savepoint: Option<Savepoint>) {
1218 if let Some(sp) = savepoint {
1219 apply_savepoint(inner, sp);
1220 }
1221}
1222
1223fn apply_savepoint(inner: &mut TxInner, sp: Savepoint) {
1224 if let Ok(mut buf) = inner.buffer.lock() {
1225 buf.truncate(sp.buffer_len);
1226 } else {
1227 inner.buffer_failed.store(true, Ordering::Release);
1228 }
1229
1230 let Some(mut graph) = sp.staged else {
1231 inner.staged = None;
1232 return;
1233 };
1234
1235 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
1239 graph.set_mutation_recorder(Some(Arc::new(BufferingRecorder::new(
1240 inner.buffer.clone(),
1241 inner.buffer_failed.clone(),
1242 )) as Arc<dyn MutationRecorder>));
1243 }
1244 inner.staged = Some(graph);
1245}
1246
1247impl Drop for Transaction<'_> {
1248 fn drop(&mut self) {
1249 if let Ok(mut inner) = self.inner.lock() {
1253 if !inner.closed {
1254 if inner.cursor_active {
1255 inner.closed = true;
1261 } else {
1262 discard_transaction_state(&mut inner);
1263 }
1264 }
1265 }
1266 }
1267}
1268
1269#[cfg(test)]
1270mod tests {
1271 use std::thread;
1272
1273 use super::*;
1274
1275 #[test]
1276 fn buffering_recorder_latches_poisoned_buffer() {
1277 let buffer = Arc::new(Mutex::new(Vec::new()));
1278 let failed = Arc::new(AtomicBool::new(false));
1279
1280 let poisoned_buffer = buffer.clone();
1281 let _ = thread::spawn(move || {
1282 let _guard = poisoned_buffer.lock().unwrap();
1283 panic!("poison mutation buffer");
1284 })
1285 .join();
1286
1287 let recorder = BufferingRecorder::new(buffer, failed.clone());
1288 recorder.record(MutationEvent::Clear);
1289
1290 assert!(failed.load(Ordering::Acquire));
1291 }
1292}