1use std::collections::BTreeMap;
2use std::sync::{Arc, Mutex, MutexGuard};
3use std::time::{Duration, Instant};
4
5use anyhow::Result;
6use thiserror::Error;
7
8use lora_analyzer::Analyzer;
9use lora_compiler::{CompiledQuery, Compiler};
10use lora_executor::{
11 classify_stream, compiled_result_columns, project_rows, ExecuteOptions, ExecutionContext,
12 Executor, LoraValue, MutableExecutionContext, MutableExecutor, MutablePullExecutor,
13 PullExecutor, QueryResult, Row, RowSource,
14};
15use lora_parser::parse_query;
16use lora_store::{InMemoryGraph, MutationEvent, MutationRecorder};
17use lora_wal::WalRecorder;
18
19use crate::error::LoraError;
20use crate::explain::{OperatorMetrics, PlanShape, ProfileMetrics, QueryPlan, QueryProfile};
21use crate::live_store::LiveStore;
22use crate::snapshot::ManagedSnapshotStore;
23use crate::stream::QueryStream;
24use crate::wal::write_scope::ensure_wal_not_poisoned;
25use lora_compiler::plan_tree_from_compiled;
26use lora_executor::{plan_result_columns, CollectorGuard, MetricsCollector};
27
28#[derive(Debug, Clone, PartialEq, Eq, Error)]
34pub enum TransactionError {
35 #[error("transaction is already closed")]
36 AlreadyClosed,
37
38 #[error("transaction has no live graph guard")]
39 NoGraphGuard,
40
41 #[error("transaction has no staged graph")]
42 NoStagedGraph,
43
44 #[error("cannot commit transaction while a streaming cursor is still active")]
45 CursorActiveCommit,
46
47 #[error("cannot start a new statement while a streaming cursor is still active")]
48 CursorActiveStatement,
49
50 #[error("cannot execute mutating query in read-only transaction")]
51 ReadOnlyMutation,
52
53 #[error("streaming write cursor requires a ReadWrite transaction")]
54 StreamingRequiresReadWrite,
55
56 #[error("read-only transaction cannot publish staged graph")]
57 ReadOnlyCommit,
58}
59
60#[derive(Debug, Clone, Copy, PartialEq, Eq)]
62pub enum TransactionMode {
63 ReadOnly,
65 ReadWrite,
67}
68
69pub(crate) enum LiveStoreGuard<'db> {
79 Read(Arc<InMemoryGraph>),
80 Write(WriteLease<'db>),
81}
82
83pub(crate) struct WriteLease<'db> {
88 pub(crate) _writer_lock: MutexGuard<'db, ()>,
90 pub(crate) store: Arc<LiveStore<InMemoryGraph>>,
92 pub(crate) snapshot: Arc<InMemoryGraph>,
95}
96
97impl LiveStoreGuard<'_> {
98 fn as_graph(&self) -> &InMemoryGraph {
99 match self {
100 Self::Read(arc) => arc,
101 Self::Write(lease) => &lease.snapshot,
102 }
103 }
104}
105
106pub(crate) struct Savepoint {
111 staged: Option<InMemoryGraph>,
112 buffer_len: usize,
113}
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
119pub(crate) enum TxStreamOutcome {
120 Exhausted,
121 Interrupted,
122}
123
124impl TxStreamOutcome {
125 fn should_restore_savepoint(self, rollback_on_drop: bool) -> bool {
126 matches!(self, Self::Interrupted) && rollback_on_drop
127 }
128}
129
130pub(crate) struct TxCursorLease {
137 handle: Arc<Mutex<TxInner>>,
138 rollback_on_drop: bool,
139 finalized: bool,
140}
141
142impl TxCursorLease {
143 pub(crate) fn new(handle: Arc<Mutex<TxInner>>, rollback_on_drop: bool) -> Self {
144 Self {
145 handle,
146 rollback_on_drop,
147 finalized: false,
148 }
149 }
150
151 pub(crate) fn finalize(&mut self, outcome: TxStreamOutcome) {
152 if self.finalized {
153 return;
154 }
155 finalize_tx_stream(&self.handle, outcome, self.rollback_on_drop);
156 self.finalized = true;
157 }
158}
159
160impl Drop for TxCursorLease {
161 fn drop(&mut self) {
162 self.finalize(TxStreamOutcome::Interrupted);
163 }
164}
165
166pub(crate) struct BufferingRecorder {
173 buffer: Arc<Mutex<Vec<MutationEvent>>>,
174}
175
176impl BufferingRecorder {
177 pub(crate) fn new(buffer: Arc<Mutex<Vec<MutationEvent>>>) -> Self {
178 Self { buffer }
179 }
180}
181
182impl MutationRecorder for BufferingRecorder {
183 fn record(&self, event: MutationEvent) {
184 if let Ok(mut buf) = self.buffer.lock() {
185 buf.push(event);
186 }
187 }
188}
189
190pub(crate) struct TxInner {
195 pub(crate) staged: Option<InMemoryGraph>,
199 pub(crate) buffer: Arc<Mutex<Vec<MutationEvent>>>,
203 pub(crate) pending_savepoint: Option<Savepoint>,
207 pub(crate) cursor_active: bool,
211 pub(crate) closed: bool,
215 pub(crate) mode: TransactionMode,
217 pub(crate) buffer_mutations: bool,
221}
222
223pub struct Transaction<'db> {
239 pub(crate) live: Option<LiveStoreGuard<'db>>,
240 pub(crate) inner: Arc<Mutex<TxInner>>,
241 pub(crate) wal: Option<Arc<WalRecorder>>,
242 pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
243 mode: TransactionMode,
244}
245
246impl<'db> Transaction<'db> {
247 pub(crate) fn new(
249 live: LiveStoreGuard<'db>,
250 wal: Option<Arc<WalRecorder>>,
251 snapshots: Option<Arc<ManagedSnapshotStore>>,
252 mode: TransactionMode,
253 ) -> Self {
254 let buffer_mutations = wal.is_some();
255 let inner = TxInner {
256 staged: None,
257 buffer: Arc::new(Mutex::new(Vec::new())),
258 pending_savepoint: None,
259 cursor_active: false,
260 closed: false,
261 mode,
262 buffer_mutations,
263 };
264 Self {
265 live: Some(live),
266 inner: Arc::new(Mutex::new(inner)),
267 wal,
268 snapshots,
269 mode,
270 }
271 }
272
273 pub fn mode(&self) -> TransactionMode {
275 self.mode
276 }
277
278 pub fn execute(
281 &mut self,
282 query: &str,
283 options: Option<ExecuteOptions>,
284 ) -> Result<QueryResult, LoraError> {
285 self.execute_with_params(query, options, BTreeMap::new())
286 }
287
288 pub fn execute_with_timeout(
290 &mut self,
291 query: &str,
292 options: Option<ExecuteOptions>,
293 timeout: Duration,
294 ) -> Result<QueryResult, LoraError> {
295 let deadline = Instant::now()
296 .checked_add(timeout)
297 .unwrap_or_else(Instant::now);
298 let rows =
299 self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
300 Ok(project_rows(rows, options.unwrap_or_default()))
301 }
302
303 pub fn execute_with_params(
305 &mut self,
306 query: &str,
307 options: Option<ExecuteOptions>,
308 params: BTreeMap<String, LoraValue>,
309 ) -> Result<QueryResult, LoraError> {
310 let rows = self.execute_rows_with_params_deadline(query, params, None)?;
311 Ok(project_rows(rows, options.unwrap_or_default()))
312 }
313
314 pub fn execute_with_params_timeout(
317 &mut self,
318 query: &str,
319 options: Option<ExecuteOptions>,
320 params: BTreeMap<String, LoraValue>,
321 timeout: Duration,
322 ) -> Result<QueryResult, LoraError> {
323 let deadline = Instant::now()
324 .checked_add(timeout)
325 .unwrap_or_else(Instant::now);
326 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
327 Ok(project_rows(rows, options.unwrap_or_default()))
328 }
329
330 pub fn execute_rows(&mut self, query: &str) -> Result<Vec<Row>, LoraError> {
333 self.execute_rows_with_params(query, BTreeMap::new())
334 }
335
336 pub fn execute_rows_with_params(
339 &mut self,
340 query: &str,
341 params: BTreeMap<String, LoraValue>,
342 ) -> Result<Vec<Row>, LoraError> {
343 Ok(self.execute_rows_with_params_deadline(query, params, None)?)
344 }
345
346 fn execute_rows_with_params_deadline(
347 &mut self,
348 query: &str,
349 params: BTreeMap<String, LoraValue>,
350 deadline: Option<Instant>,
351 ) -> Result<Vec<Row>> {
352 let compiled = self.compile_in_tx(query)?;
353 self.execute_rows_compiled_deadline(&compiled, params, deadline)
354 }
355
356 fn execute_rows_compiled_deadline(
357 &mut self,
358 compiled: &CompiledQuery,
359 params: BTreeMap<String, LoraValue>,
360 deadline: Option<Instant>,
361 ) -> Result<Vec<Row>> {
362 if self.is_read_only_unchecked() {
364 self.precheck_open_no_savepoint()?;
365 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
366 let storage = live.as_graph();
367 let executor = Executor::with_deadline(ExecutionContext { storage, params }, deadline);
368 return executor
369 .execute_compiled_rows(compiled)
370 .map_err(anyhow::Error::from);
371 }
372
373 let mut inner = self.begin_statement()?;
375 let is_mutating = classify_stream(compiled).is_mutating();
376
377 if !is_mutating {
378 return match inner.staged.as_ref() {
384 Some(staged) => {
385 let executor = Executor::with_deadline(
386 ExecutionContext {
387 storage: staged,
388 params,
389 },
390 deadline,
391 );
392 executor
393 .execute_compiled_rows(compiled)
394 .map_err(anyhow::Error::from)
395 }
396 None => {
397 drop(inner);
398 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
399 let storage = live.as_graph();
400 let executor =
401 Executor::with_deadline(ExecutionContext { storage, params }, deadline);
402 executor
403 .execute_compiled_rows(compiled)
404 .map_err(anyhow::Error::from)
405 }
406 };
407 }
408
409 let clone_savepoint_graph = inner.staged.is_some();
413 self.ensure_staged_locked(&mut inner)?;
414 let savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
415
416 let exec_result: ExecResultRows = {
417 let staged = inner.staged_mut()?;
418 let mut executor = MutableExecutor::with_deadline(
419 MutableExecutionContext {
420 storage: staged,
421 params,
422 },
423 deadline,
424 );
425 executor
426 .execute_compiled_rows(compiled)
427 .map_err(anyhow::Error::from)
428 };
429
430 match exec_result {
431 Ok(rows) => Ok(rows),
432 Err(err) => {
433 restore_savepoint(&mut inner, savepoint);
434 Err(err)
435 }
436 }
437 }
438
439 pub(crate) fn open_streaming_compiled_autocommit(
468 &mut self,
469 compiled: Arc<CompiledQuery>,
470 params: BTreeMap<String, LoraValue>,
471 ) -> Result<Box<dyn RowSource + 'static>> {
472 if self.is_read_only_unchecked() {
473 return Err(TransactionError::StreamingRequiresReadWrite.into());
474 }
475
476 let mut inner = self.begin_statement()?;
477 self.ensure_staged_locked(&mut inner)?;
478 inner.activate_cursor();
479
480 let staged_ptr: *mut InMemoryGraph = inner
484 .staged
485 .as_mut()
486 .expect("ensure_staged_locked guarantees Some")
487 as *mut _;
488 drop(inner);
489
490 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
496 let compiled_static: &'static CompiledQuery =
497 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
498
499 let cursor = MutablePullExecutor::new(storage_static, params)
503 .open_compiled(compiled_static)
504 .map_err(|e| {
505 if let Ok(mut inner) = self.inner.lock() {
509 discard_transaction_state(&mut inner);
510 }
511 self.live.take();
512 anyhow::Error::from(e)
513 })?;
514
515 Ok(Box::new(StreamingCursorWithArc {
521 cursor,
522 _compiled: compiled,
523 }))
524 }
525
526 fn compile_in_tx(&self, query: &str) -> Result<CompiledQuery> {
532 let document = parse_query(query)?;
533 let resolved = {
534 let inner = self.lock_inner_unchecked();
535 if let Some(staged) = &inner.staged {
536 let mut analyzer = Analyzer::new(staged);
537 analyzer.analyze(&document)?
538 } else {
539 drop(inner);
540 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
541 let mut analyzer = Analyzer::new(live.as_graph());
542 analyzer.analyze(&document)?
543 }
544 };
545 Ok(Compiler::compile(&resolved))
546 }
547
548 fn ensure_staged_locked(&self, inner: &mut MutexGuard<'_, TxInner>) -> Result<()> {
552 if inner.staged.is_some() {
553 return Ok(());
554 }
555 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
556 let mut staged: InMemoryGraph = live.as_graph().clone();
557 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
558 staged.set_mutation_recorder(Some(
559 Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>,
560 ));
561 }
562 inner.staged = Some(staged);
563 Ok(())
564 }
565
566 pub fn explain(
571 &self,
572 query: &str,
573 _params: Option<BTreeMap<String, LoraValue>>,
574 ) -> Result<QueryPlan, LoraError> {
575 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
576 let tree = plan_tree_from_compiled(&compiled);
577 let shape: PlanShape = classify_stream(&compiled).into();
578 let result_columns = plan_result_columns(&compiled.physical);
579 Ok(QueryPlan {
580 query: query.to_string(),
581 tree,
582 shape,
583 result_columns,
584 })
585 }
586
587 pub fn profile(
595 &mut self,
596 query: &str,
597 params: Option<BTreeMap<String, LoraValue>>,
598 ) -> Result<QueryProfile, LoraError> {
599 let params = params.unwrap_or_default();
600 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
601 let tree = plan_tree_from_compiled(&compiled);
602 let shape: PlanShape = classify_stream(&compiled).into();
603 let result_columns = plan_result_columns(&compiled.physical);
604 let plan = QueryPlan {
605 query: query.to_string(),
606 tree,
607 shape,
608 result_columns,
609 };
610
611 let collector = Arc::new(MetricsCollector::new());
612 let _guard = CollectorGuard::install(collector.clone());
613
614 let started = Instant::now();
615 let rows = self
616 .execute_rows_compiled_deadline(&compiled, params, None)
617 .map_err(LoraError::from_anyhow)?;
618 let total_elapsed_ns = started.elapsed().as_nanos() as u64;
619
620 drop(_guard);
621 let per_operator = collector
622 .snapshot()
623 .into_iter()
624 .map(|(id, op)| {
625 (
626 id,
627 OperatorMetrics {
628 rows: op.rows,
629 elapsed_ns: op.elapsed_ns,
630 next_calls: op.next_calls,
631 db_hits: 0,
632 },
633 )
634 })
635 .collect();
636
637 let metrics = ProfileMetrics {
638 total_elapsed_ns,
639 total_rows: rows.len() as u64,
640 mutated: shape.is_mutating(),
641 per_operator,
642 };
643
644 Ok(QueryProfile { plan, metrics })
645 }
646
647 pub fn stream(&mut self, query: &str) -> Result<QueryStream<'static>, LoraError> {
649 self.stream_with_params(query, BTreeMap::new())
650 }
651
652 pub fn stream_with_params(
655 &mut self,
656 query: &str,
657 params: BTreeMap<String, LoraValue>,
658 ) -> Result<QueryStream<'static>, LoraError> {
659 let compiled = Arc::new(self.compile_in_tx(query)?);
660 let columns = compiled_result_columns(&compiled);
661 Ok(self.stream_compiled(compiled, columns, params)?)
662 }
663
664 pub(crate) fn stream_compiled(
668 &mut self,
669 compiled: Arc<CompiledQuery>,
670 columns: Vec<String>,
671 params: BTreeMap<String, LoraValue>,
672 ) -> Result<QueryStream<'static>> {
673 let mut inner = self.begin_statement()?;
674 let is_mutating = classify_stream(&compiled).is_mutating();
675 if matches!(inner.mode, TransactionMode::ReadOnly) && is_mutating {
676 return Err(TransactionError::ReadOnlyMutation.into());
677 }
678
679 let clone_savepoint_graph = inner.staged.is_some();
684 self.ensure_staged_locked(&mut inner)?;
685 inner.activate_cursor();
686
687 let rollback_on_drop = is_mutating;
688 if rollback_on_drop {
689 inner.pending_savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
690 } else {
691 inner.pending_savepoint = None;
692 }
693
694 let staged_ptr: *mut InMemoryGraph = inner
695 .staged
696 .as_mut()
697 .expect("ensure_staged_locked guarantees Some")
698 as *mut _;
699 drop(inner);
700
701 let compiled_static: &'static CompiledQuery =
702 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
703 let cursor: Result<Box<dyn RowSource + 'static>> = if is_mutating {
704 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
705 MutablePullExecutor::new(storage_static, params)
706 .open_compiled(compiled_static)
707 .map(|cursor| {
708 Box::new(StreamingCursorWithArc {
709 cursor,
710 _compiled: compiled.clone(),
711 }) as Box<dyn RowSource + 'static>
712 })
713 .map_err(anyhow::Error::from)
714 } else {
715 let storage_static: &'static InMemoryGraph = unsafe { &*staged_ptr };
716 PullExecutor::new(storage_static, params)
717 .open_compiled(compiled_static)
718 .map(|cursor| {
719 Box::new(StreamingCursorWithArc {
720 cursor,
721 _compiled: compiled.clone(),
722 }) as Box<dyn RowSource + 'static>
723 })
724 .map_err(anyhow::Error::from)
725 };
726
727 match cursor {
728 Ok(cursor) => Ok(QueryStream::for_tx_cursor(
729 cursor,
730 columns,
731 TxCursorLease::new(self.inner.clone(), rollback_on_drop),
732 )),
733 Err(err) => {
734 finalize_tx_stream(&self.inner, TxStreamOutcome::Interrupted, rollback_on_drop);
735 Err(err)
736 }
737 }
738 }
739
740 pub fn commit(mut self) -> Result<(), LoraError> {
747 let CommitState {
748 staged,
749 buffer_events,
750 mode,
751 } = self.take_commit_state()?;
752
753 let wrote_wal_commit = self.replay_commit_wal(mode, buffer_events)?;
754 self.publish_staged_graph(mode, staged, wrote_wal_commit)?;
755
756 self.live.take();
757 Ok(())
758 }
759
760 fn take_commit_state(&self) -> Result<CommitState> {
761 let mut inner = self.inner.lock().unwrap();
762 if inner.cursor_active {
763 return Err(TransactionError::CursorActiveCommit.into());
764 }
765 if inner.closed {
766 return Err(TransactionError::AlreadyClosed.into());
767 }
768
769 let mode = inner.mode;
770 let staged = inner.staged.take();
774 let buffer_events = std::mem::take(&mut *inner.buffer.lock().unwrap());
775 inner.closed = true;
776
777 Ok(CommitState {
778 staged,
779 buffer_events,
780 mode,
781 })
782 }
783
784 fn replay_commit_wal(
785 &self,
786 mode: TransactionMode,
787 buffer_events: Vec<MutationEvent>,
788 ) -> Result<bool> {
789 let Some(rec) = &self.wal else {
790 return Ok(false);
791 };
792
793 if !matches!(mode, TransactionMode::ReadWrite) {
794 ensure_wal_not_poisoned(rec)?;
795 return Ok(false);
796 }
797
798 Ok(rec.commit_events(buffer_events)?.wrote())
799 }
800
801 fn publish_staged_graph(
802 &mut self,
803 mode: TransactionMode,
804 staged: Option<InMemoryGraph>,
805 wrote_wal_commit: bool,
806 ) -> Result<()> {
807 if !matches!(mode, TransactionMode::ReadWrite) {
808 return Ok(());
809 }
810
811 let Some(mut staged) = staged else {
812 return Ok(());
813 };
814
815 staged.set_mutation_recorder(None);
820 let wal = self.wal.clone();
821 if let Some(rec) = &wal {
822 staged.set_mutation_recorder(Some(rec.clone() as Arc<dyn MutationRecorder>));
823 }
824
825 let live = self.live.as_mut().ok_or(TransactionError::NoGraphGuard)?;
826 let lease = match live {
827 LiveStoreGuard::Write(lease) => lease,
828 LiveStoreGuard::Read(_) => {
829 return Err(TransactionError::ReadOnlyCommit.into());
830 }
831 };
832
833 if wrote_wal_commit {
834 if let (Some(snapshots), Some(rec)) = (&self.snapshots, wal.as_ref()) {
835 snapshots.observe_commit(&staged, rec)?;
836 }
837 }
838
839 lease.store.store(Arc::new(staged));
843
844 Ok(())
845 }
846
847 pub fn rollback(mut self) -> Result<(), LoraError> {
850 let mut inner = self.inner.lock().unwrap();
851 if inner.closed {
852 return Err(TransactionError::AlreadyClosed.into());
853 }
854 discard_transaction_state(&mut inner);
855 drop(inner);
856 self.live.take();
857 Ok(())
858 }
859
860 fn begin_statement(&self) -> Result<MutexGuard<'_, TxInner>> {
866 let inner = self.inner.lock().unwrap();
867 if inner.closed {
868 return Err(TransactionError::AlreadyClosed.into());
869 }
870 if inner.cursor_active {
871 return Err(TransactionError::CursorActiveStatement.into());
872 }
873 Ok(inner)
874 }
875
876 fn precheck_open_no_savepoint(&self) -> Result<()> {
880 let inner = self.inner.lock().unwrap();
881 if inner.closed {
882 return Err(TransactionError::AlreadyClosed.into());
883 }
884 if inner.cursor_active {
885 return Err(TransactionError::CursorActiveStatement.into());
886 }
887 Ok(())
888 }
889
890 fn is_read_only_unchecked(&self) -> bool {
894 matches!(self.mode, TransactionMode::ReadOnly)
895 }
896
897 fn lock_inner_unchecked(&self) -> MutexGuard<'_, TxInner> {
898 self.inner
899 .lock()
900 .unwrap_or_else(|poisoned| poisoned.into_inner())
901 }
902
903 pub(crate) fn release_streaming_cursor(&self) {
904 if let Ok(mut inner) = self.inner.lock() {
905 inner.release_cursor();
906 }
907 }
908}
909
910type ExecResultRows = Result<Vec<Row>>;
911
912struct CommitState {
913 staged: Option<InMemoryGraph>,
914 buffer_events: Vec<MutationEvent>,
915 mode: TransactionMode,
916}
917
918impl TxInner {
919 fn staged_mut(&mut self) -> Result<&mut InMemoryGraph> {
920 self.staged
921 .as_mut()
922 .ok_or(TransactionError::NoStagedGraph.into())
923 }
924
925 fn activate_cursor(&mut self) {
926 self.cursor_active = true;
927 }
928
929 fn release_cursor(&mut self) {
930 self.cursor_active = false;
931 }
932
933 fn clear_pending_savepoint(&mut self) {
934 self.pending_savepoint = None;
935 }
936
937 fn restore_pending_savepoint(&mut self) {
938 if let Some(sp) = self.pending_savepoint.take() {
939 apply_savepoint(self, sp);
940 }
941 }
942
943 fn finalize_stream(&mut self, outcome: TxStreamOutcome, rollback_on_drop: bool) {
944 self.release_cursor();
945
946 if self.closed {
947 discard_transaction_state(self);
948 return;
949 }
950
951 if outcome.should_restore_savepoint(rollback_on_drop) {
952 self.restore_pending_savepoint();
953 } else {
954 self.clear_pending_savepoint();
955 }
956 }
957}
958
959struct StreamingCursorWithArc {
964 cursor: Box<dyn RowSource + 'static>,
965 _compiled: Arc<CompiledQuery>,
966}
967
968impl RowSource for StreamingCursorWithArc {
969 fn next_row(&mut self) -> lora_executor::ExecResult<Option<Row>> {
970 self.cursor.next_row()
971 }
972}
973
974fn finalize_tx_stream(
975 handle: &Arc<Mutex<TxInner>>,
976 outcome: TxStreamOutcome,
977 rollback_on_drop: bool,
978) {
979 if let Ok(mut inner) = handle.lock() {
980 inner.finalize_stream(outcome, rollback_on_drop);
981 }
982}
983
984fn discard_transaction_state(inner: &mut TxInner) {
985 inner.clear_pending_savepoint();
987 inner.release_cursor();
988 inner.staged = None;
989 if let Ok(mut buf) = inner.buffer.lock() {
990 buf.clear();
991 }
992 inner.closed = true;
993}
994
995fn take_savepoint(inner: &TxInner, clone_staged: bool) -> Savepoint {
996 let buffer_len = inner.buffer.lock().ok().map(|b| b.len()).unwrap_or(0);
997 Savepoint {
998 staged: if clone_staged {
999 inner.staged.as_ref().cloned()
1000 } else {
1001 None
1002 },
1003 buffer_len,
1004 }
1005}
1006
1007fn restore_savepoint(inner: &mut TxInner, savepoint: Option<Savepoint>) {
1008 if let Some(sp) = savepoint {
1009 apply_savepoint(inner, sp);
1010 }
1011}
1012
1013fn apply_savepoint(inner: &mut TxInner, sp: Savepoint) {
1014 if let Ok(mut buf) = inner.buffer.lock() {
1015 buf.truncate(sp.buffer_len);
1016 }
1017
1018 let Some(mut graph) = sp.staged else {
1019 inner.staged = None;
1020 return;
1021 };
1022
1023 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
1027 graph.set_mutation_recorder(Some(
1028 Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>
1029 ));
1030 }
1031 inner.staged = Some(graph);
1032}
1033
1034impl Drop for Transaction<'_> {
1035 fn drop(&mut self) {
1036 if let Ok(mut inner) = self.inner.lock() {
1040 if !inner.closed {
1041 if inner.cursor_active {
1042 inner.closed = true;
1048 } else {
1049 discard_transaction_state(&mut inner);
1050 }
1051 }
1052 }
1053 }
1054}