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 {
178 buffer: Arc<Mutex<Vec<MutationEvent>>>,
179}
180
181impl BufferingRecorder {
182 pub(crate) fn new(buffer: Arc<Mutex<Vec<MutationEvent>>>) -> Self {
183 Self { buffer }
184 }
185}
186
187impl MutationRecorder for BufferingRecorder {
188 fn record(&self, event: MutationEvent) {
189 if let Ok(mut buf) = self.buffer.lock() {
190 buf.push(event);
191 }
192 }
193}
194
195pub(crate) struct TxInner {
200 pub(crate) staged: Option<InMemoryGraph>,
204 pub(crate) buffer: Arc<Mutex<Vec<MutationEvent>>>,
208 pub(crate) pending_savepoint: Option<Savepoint>,
212 pub(crate) cursor_active: bool,
216 pub(crate) closed: bool,
220 pub(crate) mode: TransactionMode,
222 pub(crate) buffer_mutations: bool,
226}
227
228pub struct Transaction<'db> {
244 pub(crate) live: Option<LiveStoreGuard<'db>>,
245 pub(crate) inner: Arc<Mutex<TxInner>>,
246 pub(crate) wal: Option<Arc<WalRecorder>>,
247 pub(crate) snapshots: Option<Arc<ManagedSnapshotStore>>,
248 mode: TransactionMode,
249}
250
251impl<'db> Transaction<'db> {
252 pub(crate) fn new(
254 live: LiveStoreGuard<'db>,
255 wal: Option<Arc<WalRecorder>>,
256 snapshots: Option<Arc<ManagedSnapshotStore>>,
257 mode: TransactionMode,
258 ) -> Self {
259 let buffer_mutations = wal.is_some();
260 let inner = TxInner {
261 staged: None,
262 buffer: Arc::new(Mutex::new(Vec::new())),
263 pending_savepoint: None,
264 cursor_active: false,
265 closed: false,
266 mode,
267 buffer_mutations,
268 };
269 Self {
270 live: Some(live),
271 inner: Arc::new(Mutex::new(inner)),
272 wal,
273 snapshots,
274 mode,
275 }
276 }
277
278 pub fn mode(&self) -> TransactionMode {
280 self.mode
281 }
282
283 pub fn execute(
286 &mut self,
287 query: &str,
288 options: Option<ExecuteOptions>,
289 ) -> Result<QueryResult, LoraError> {
290 self.execute_with_params(query, options, BTreeMap::new())
291 }
292
293 pub fn execute_with_timeout(
295 &mut self,
296 query: &str,
297 options: Option<ExecuteOptions>,
298 timeout: Duration,
299 ) -> Result<QueryResult, LoraError> {
300 let deadline = Instant::now()
301 .checked_add(timeout)
302 .unwrap_or_else(Instant::now);
303 let rows =
304 self.execute_rows_with_params_deadline(query, BTreeMap::new(), Some(deadline))?;
305 Ok(project_rows(rows, options.unwrap_or_default()))
306 }
307
308 pub fn execute_with_params(
310 &mut self,
311 query: &str,
312 options: Option<ExecuteOptions>,
313 params: BTreeMap<String, LoraValue>,
314 ) -> Result<QueryResult, LoraError> {
315 let rows = self.execute_rows_with_params_deadline(query, params, None)?;
316 Ok(project_rows(rows, options.unwrap_or_default()))
317 }
318
319 pub fn execute_with_params_timeout(
322 &mut self,
323 query: &str,
324 options: Option<ExecuteOptions>,
325 params: BTreeMap<String, LoraValue>,
326 timeout: Duration,
327 ) -> Result<QueryResult, LoraError> {
328 let deadline = Instant::now()
329 .checked_add(timeout)
330 .unwrap_or_else(Instant::now);
331 let rows = self.execute_rows_with_params_deadline(query, params, Some(deadline))?;
332 Ok(project_rows(rows, options.unwrap_or_default()))
333 }
334
335 pub fn execute_rows(&mut self, query: &str) -> Result<Vec<Row>, LoraError> {
338 self.execute_rows_with_params(query, BTreeMap::new())
339 }
340
341 pub fn execute_rows_with_params(
344 &mut self,
345 query: &str,
346 params: BTreeMap<String, LoraValue>,
347 ) -> Result<Vec<Row>, LoraError> {
348 Ok(self.execute_rows_with_params_deadline(query, params, None)?)
349 }
350
351 fn execute_rows_with_params_deadline(
352 &mut self,
353 query: &str,
354 params: BTreeMap<String, LoraValue>,
355 deadline: Option<Instant>,
356 ) -> Result<Vec<Row>> {
357 let compiled = self.compile_in_tx(query)?;
358 self.execute_rows_compiled_deadline(&compiled, params, deadline)
359 }
360
361 fn execute_rows_compiled_deadline(
362 &mut self,
363 compiled: &CompiledQuery,
364 params: BTreeMap<String, LoraValue>,
365 deadline: Option<Instant>,
366 ) -> Result<Vec<Row>> {
367 if self.is_read_only_unchecked() {
369 self.precheck_open_no_savepoint()?;
370 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
371 let storage = live.as_graph();
372 let executor = Executor::with_deadline(ExecutionContext { storage, params }, deadline);
373 return executor
374 .execute_compiled_rows(compiled)
375 .map_err(anyhow::Error::from);
376 }
377
378 let mut inner = self.begin_statement()?;
380 let is_mutating = classify_stream(compiled).is_mutating();
381
382 if !is_mutating {
383 return match inner.staged.as_ref() {
389 Some(staged) => {
390 let executor = Executor::with_deadline(
391 ExecutionContext {
392 storage: staged,
393 params,
394 },
395 deadline,
396 );
397 executor
398 .execute_compiled_rows(compiled)
399 .map_err(anyhow::Error::from)
400 }
401 None => {
402 drop(inner);
403 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
404 let storage = live.as_graph();
405 let executor =
406 Executor::with_deadline(ExecutionContext { storage, params }, deadline);
407 executor
408 .execute_compiled_rows(compiled)
409 .map_err(anyhow::Error::from)
410 }
411 };
412 }
413
414 let clone_savepoint_graph = inner.staged.is_some();
418 self.ensure_staged_locked(&mut inner)?;
419 let savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
420
421 let exec_result: ExecResultRows = {
422 let staged = inner.staged_mut()?;
423 let mut executor = MutableExecutor::with_deadline(
424 MutableExecutionContext {
425 storage: staged,
426 params,
427 },
428 deadline,
429 );
430 executor
431 .execute_compiled_rows(compiled)
432 .map_err(anyhow::Error::from)
433 };
434
435 match exec_result {
436 Ok(rows) => Ok(rows),
437 Err(err) => {
438 restore_savepoint(&mut inner, savepoint);
439 Err(err)
440 }
441 }
442 }
443
444 pub(crate) fn open_streaming_compiled_autocommit(
473 &mut self,
474 compiled: Arc<CompiledQuery>,
475 params: BTreeMap<String, LoraValue>,
476 ) -> Result<Box<dyn RowSource + 'static>> {
477 if self.is_read_only_unchecked() {
478 return Err(TransactionError::StreamingRequiresReadWrite.into());
479 }
480
481 let mut inner = self.begin_statement()?;
482 self.ensure_staged_locked(&mut inner)?;
483 inner.activate_cursor();
484
485 let staged_ptr: *mut InMemoryGraph = inner
489 .staged
490 .as_mut()
491 .expect("ensure_staged_locked guarantees Some")
492 as *mut _;
493 drop(inner);
494
495 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
501 let compiled_static: &'static CompiledQuery =
502 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
503
504 let cursor = MutablePullExecutor::new(storage_static, params)
508 .open_compiled(compiled_static)
509 .map_err(|e| {
510 if let Ok(mut inner) = self.inner.lock() {
514 discard_transaction_state(&mut inner);
515 }
516 self.live.take();
517 anyhow::Error::from(e)
518 })?;
519
520 Ok(Box::new(StreamingCursorWithArc {
526 cursor,
527 _compiled: compiled,
528 }))
529 }
530
531 fn compile_in_tx(&self, query: &str) -> Result<CompiledQuery> {
537 let document = parse_query(query)?;
538 let resolved = {
539 let inner = self.lock_inner_unchecked();
540 if let Some(staged) = &inner.staged {
541 let mut analyzer = Analyzer::new(staged);
542 analyzer.analyze(&document)?
543 } else {
544 drop(inner);
545 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
546 let mut analyzer = Analyzer::new(live.as_graph());
547 analyzer.analyze(&document)?
548 }
549 };
550 Ok(Compiler::compile(&resolved))
551 }
552
553 fn ensure_staged_locked(&self, inner: &mut MutexGuard<'_, TxInner>) -> Result<()> {
557 if inner.staged.is_some() {
558 return Ok(());
559 }
560 let live = self.live.as_ref().ok_or(TransactionError::NoGraphGuard)?;
561 let mut staged: InMemoryGraph = live.as_graph().clone();
562 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
563 staged.set_mutation_recorder(Some(
564 Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>,
565 ));
566 }
567 inner.staged = Some(staged);
568 Ok(())
569 }
570
571 pub fn explain(
576 &self,
577 query: &str,
578 _params: Option<BTreeMap<String, LoraValue>>,
579 ) -> Result<QueryPlan, LoraError> {
580 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
581 let tree = plan_tree_from_compiled(&compiled);
582 let shape: PlanShape = classify_stream(&compiled).into();
583 let result_columns = plan_result_columns(&compiled.physical);
584 Ok(QueryPlan {
585 query: query.to_string(),
586 tree,
587 shape,
588 result_columns,
589 })
590 }
591
592 pub fn profile(
600 &mut self,
601 query: &str,
602 params: Option<BTreeMap<String, LoraValue>>,
603 ) -> Result<QueryProfile, LoraError> {
604 let params = params.unwrap_or_default();
605 let compiled = self.compile_in_tx(query).map_err(LoraError::from_anyhow)?;
606 let tree = plan_tree_from_compiled(&compiled);
607 let shape: PlanShape = classify_stream(&compiled).into();
608 let result_columns = plan_result_columns(&compiled.physical);
609 let plan = QueryPlan {
610 query: query.to_string(),
611 tree,
612 shape,
613 result_columns,
614 };
615
616 let collector = Arc::new(MetricsCollector::new());
617 let _guard = CollectorGuard::install(collector.clone());
618
619 let started = Instant::now();
620 let rows = self
621 .execute_rows_compiled_deadline(&compiled, params, None)
622 .map_err(LoraError::from_anyhow)?;
623 let total_elapsed_ns = started.elapsed().as_nanos() as u64;
624
625 drop(_guard);
626 let per_operator = collector
627 .snapshot()
628 .into_iter()
629 .map(|(id, op)| {
630 (
631 id,
632 OperatorMetrics {
633 rows: op.rows,
634 elapsed_ns: op.elapsed_ns,
635 next_calls: op.next_calls,
636 db_hits: 0,
637 },
638 )
639 })
640 .collect();
641
642 let metrics = ProfileMetrics {
643 total_elapsed_ns,
644 total_rows: rows.len() as u64,
645 mutated: shape.is_mutating(),
646 per_operator,
647 };
648
649 Ok(QueryProfile { plan, metrics })
650 }
651
652 pub fn stream(&mut self, query: &str) -> Result<QueryStream<'static>, LoraError> {
654 self.stream_with_params(query, BTreeMap::new())
655 }
656
657 pub fn stream_with_params(
660 &mut self,
661 query: &str,
662 params: BTreeMap<String, LoraValue>,
663 ) -> Result<QueryStream<'static>, LoraError> {
664 let compiled = Arc::new(self.compile_in_tx(query)?);
665 let columns = compiled_result_columns(&compiled);
666 Ok(self.stream_compiled(compiled, columns, params)?)
667 }
668
669 pub(crate) fn stream_compiled(
673 &mut self,
674 compiled: Arc<CompiledQuery>,
675 columns: Vec<String>,
676 params: BTreeMap<String, LoraValue>,
677 ) -> Result<QueryStream<'static>> {
678 let mut inner = self.begin_statement()?;
679 let is_mutating = classify_stream(&compiled).is_mutating();
680 if matches!(inner.mode, TransactionMode::ReadOnly) && is_mutating {
681 return Err(TransactionError::ReadOnlyMutation.into());
682 }
683
684 let clone_savepoint_graph = inner.staged.is_some();
689 self.ensure_staged_locked(&mut inner)?;
690 inner.activate_cursor();
691
692 let rollback_on_drop = is_mutating;
693 if rollback_on_drop {
694 inner.pending_savepoint = Some(take_savepoint(&inner, clone_savepoint_graph));
695 } else {
696 inner.pending_savepoint = None;
697 }
698
699 let staged_ptr: *mut InMemoryGraph = inner
700 .staged
701 .as_mut()
702 .expect("ensure_staged_locked guarantees Some")
703 as *mut _;
704 drop(inner);
705
706 let compiled_static: &'static CompiledQuery =
707 unsafe { std::mem::transmute::<&CompiledQuery, _>(compiled.as_ref()) };
708 let cursor: Result<Box<dyn RowSource + 'static>> = if is_mutating {
709 let storage_static: &'static mut InMemoryGraph = unsafe { &mut *staged_ptr };
710 MutablePullExecutor::new(storage_static, params)
711 .open_compiled(compiled_static)
712 .map(|cursor| {
713 Box::new(StreamingCursorWithArc {
714 cursor,
715 _compiled: compiled.clone(),
716 }) as Box<dyn RowSource + 'static>
717 })
718 .map_err(anyhow::Error::from)
719 } else {
720 let storage_static: &'static InMemoryGraph = unsafe { &*staged_ptr };
721 PullExecutor::new(storage_static, params)
722 .open_compiled(compiled_static)
723 .map(|cursor| {
724 Box::new(StreamingCursorWithArc {
725 cursor,
726 _compiled: compiled.clone(),
727 }) as Box<dyn RowSource + 'static>
728 })
729 .map_err(anyhow::Error::from)
730 };
731
732 match cursor {
733 Ok(cursor) => Ok(QueryStream::for_tx_cursor(
734 cursor,
735 columns,
736 TxCursorLease::new(self.inner.clone(), rollback_on_drop),
737 )),
738 Err(err) => {
739 finalize_tx_stream(&self.inner, TxStreamOutcome::Interrupted, rollback_on_drop);
740 Err(err)
741 }
742 }
743 }
744
745 pub fn commit(mut self) -> Result<(), LoraError> {
752 let CommitState {
753 staged,
754 buffer_events,
755 mode,
756 } = self.take_commit_state()?;
757
758 let wrote_wal_commit = self.replay_commit_wal(mode, buffer_events)?;
759 self.publish_staged_graph(mode, staged, wrote_wal_commit)?;
760
761 self.live.take();
762 Ok(())
763 }
764
765 fn take_commit_state(&self) -> Result<CommitState> {
766 let mut inner = self.inner.lock().unwrap();
767 if inner.cursor_active {
768 return Err(TransactionError::CursorActiveCommit.into());
769 }
770 if inner.closed {
771 return Err(TransactionError::AlreadyClosed.into());
772 }
773
774 let mode = inner.mode;
775 let staged = inner.staged.take();
779 let buffer_events = std::mem::take(&mut *inner.buffer.lock().unwrap());
780 inner.closed = true;
781
782 Ok(CommitState {
783 staged,
784 buffer_events,
785 mode,
786 })
787 }
788
789 fn replay_commit_wal(
790 &self,
791 mode: TransactionMode,
792 buffer_events: Vec<MutationEvent>,
793 ) -> Result<bool> {
794 let Some(rec) = &self.wal else {
795 return Ok(false);
796 };
797
798 if !matches!(mode, TransactionMode::ReadWrite) {
799 ensure_wal_not_poisoned(rec)?;
800 return Ok(false);
801 }
802
803 Ok(rec.commit_events(buffer_events)?.wrote())
804 }
805
806 fn publish_staged_graph(
807 &mut self,
808 mode: TransactionMode,
809 staged: Option<InMemoryGraph>,
810 wrote_wal_commit: bool,
811 ) -> Result<()> {
812 if !matches!(mode, TransactionMode::ReadWrite) {
813 return Ok(());
814 }
815
816 let Some(mut staged) = staged else {
817 return Ok(());
818 };
819
820 staged.set_mutation_recorder(None);
825 let wal = self.wal.clone();
826 if let Some(rec) = &wal {
827 staged.set_mutation_recorder(Some(rec.clone() as Arc<dyn MutationRecorder>));
828 }
829
830 let live = self.live.as_mut().ok_or(TransactionError::NoGraphGuard)?;
831 let lease = match live {
832 LiveStoreGuard::Write(lease) => lease,
833 LiveStoreGuard::Read(_) => {
834 return Err(TransactionError::ReadOnlyCommit.into());
835 }
836 };
837
838 if wrote_wal_commit {
839 if let (Some(snapshots), Some(rec)) = (&self.snapshots, wal.as_ref()) {
840 snapshots.observe_commit(&staged, rec)?;
841 }
842 }
843
844 lease.store.store(Arc::new(staged));
848
849 Ok(())
850 }
851
852 pub fn rollback(mut self) -> Result<(), LoraError> {
855 let mut inner = self.inner.lock().unwrap();
856 if inner.closed {
857 return Err(TransactionError::AlreadyClosed.into());
858 }
859 discard_transaction_state(&mut inner);
860 drop(inner);
861 self.live.take();
862 Ok(())
863 }
864
865 fn begin_statement(&self) -> Result<MutexGuard<'_, TxInner>> {
871 let inner = self.inner.lock().unwrap();
872 if inner.closed {
873 return Err(TransactionError::AlreadyClosed.into());
874 }
875 if inner.cursor_active {
876 return Err(TransactionError::CursorActiveStatement.into());
877 }
878 Ok(inner)
879 }
880
881 fn precheck_open_no_savepoint(&self) -> Result<()> {
885 let inner = self.inner.lock().unwrap();
886 if inner.closed {
887 return Err(TransactionError::AlreadyClosed.into());
888 }
889 if inner.cursor_active {
890 return Err(TransactionError::CursorActiveStatement.into());
891 }
892 Ok(())
893 }
894
895 fn is_read_only_unchecked(&self) -> bool {
899 matches!(self.mode, TransactionMode::ReadOnly)
900 }
901
902 fn lock_inner_unchecked(&self) -> MutexGuard<'_, TxInner> {
903 self.inner
904 .lock()
905 .unwrap_or_else(|poisoned| poisoned.into_inner())
906 }
907
908 pub(crate) fn release_streaming_cursor(&self) {
909 if let Ok(mut inner) = self.inner.lock() {
910 inner.release_cursor();
911 }
912 }
913}
914
915type ExecResultRows = Result<Vec<Row>>;
916
917struct CommitState {
918 staged: Option<InMemoryGraph>,
919 buffer_events: Vec<MutationEvent>,
920 mode: TransactionMode,
921}
922
923impl TxInner {
924 fn staged_mut(&mut self) -> Result<&mut InMemoryGraph> {
925 self.staged
926 .as_mut()
927 .ok_or(TransactionError::NoStagedGraph.into())
928 }
929
930 fn activate_cursor(&mut self) {
931 self.cursor_active = true;
932 }
933
934 fn release_cursor(&mut self) {
935 self.cursor_active = false;
936 }
937
938 fn clear_pending_savepoint(&mut self) {
939 self.pending_savepoint = None;
940 }
941
942 fn restore_pending_savepoint(&mut self) {
943 if let Some(sp) = self.pending_savepoint.take() {
944 apply_savepoint(self, sp);
945 }
946 }
947
948 fn finalize_stream(&mut self, outcome: TxStreamOutcome, rollback_on_drop: bool) {
949 self.release_cursor();
950
951 if self.closed {
952 discard_transaction_state(self);
953 return;
954 }
955
956 if outcome.should_restore_savepoint(rollback_on_drop) {
957 self.restore_pending_savepoint();
958 } else {
959 self.clear_pending_savepoint();
960 }
961 }
962}
963
964struct StreamingCursorWithArc {
969 cursor: Box<dyn RowSource + 'static>,
970 _compiled: Arc<CompiledQuery>,
971}
972
973impl RowSource for StreamingCursorWithArc {
974 fn next_row(&mut self) -> lora_executor::ExecResult<Option<Row>> {
975 self.cursor.next_row()
976 }
977}
978
979fn finalize_tx_stream(
980 handle: &Arc<Mutex<TxInner>>,
981 outcome: TxStreamOutcome,
982 rollback_on_drop: bool,
983) {
984 if let Ok(mut inner) = handle.lock() {
985 inner.finalize_stream(outcome, rollback_on_drop);
986 }
987}
988
989fn discard_transaction_state(inner: &mut TxInner) {
990 inner.clear_pending_savepoint();
992 inner.release_cursor();
993 inner.staged = None;
994 if let Ok(mut buf) = inner.buffer.lock() {
995 buf.clear();
996 }
997 inner.closed = true;
998}
999
1000fn take_savepoint(inner: &TxInner, clone_staged: bool) -> Savepoint {
1001 let buffer_len = inner.buffer.lock().ok().map(|b| b.len()).unwrap_or(0);
1002 Savepoint {
1003 staged: if clone_staged {
1004 inner.staged.as_ref().cloned()
1005 } else {
1006 None
1007 },
1008 buffer_len,
1009 }
1010}
1011
1012fn restore_savepoint(inner: &mut TxInner, savepoint: Option<Savepoint>) {
1013 if let Some(sp) = savepoint {
1014 apply_savepoint(inner, sp);
1015 }
1016}
1017
1018fn apply_savepoint(inner: &mut TxInner, sp: Savepoint) {
1019 if let Ok(mut buf) = inner.buffer.lock() {
1020 buf.truncate(sp.buffer_len);
1021 }
1022
1023 let Some(mut graph) = sp.staged else {
1024 inner.staged = None;
1025 return;
1026 };
1027
1028 if matches!(inner.mode, TransactionMode::ReadWrite) && inner.buffer_mutations {
1032 graph.set_mutation_recorder(Some(
1033 Arc::new(BufferingRecorder::new(inner.buffer.clone())) as Arc<dyn MutationRecorder>
1034 ));
1035 }
1036 inner.staged = Some(graph);
1037}
1038
1039impl Drop for Transaction<'_> {
1040 fn drop(&mut self) {
1041 if let Ok(mut inner) = self.inner.lock() {
1045 if !inner.closed {
1046 if inner.cursor_active {
1047 inner.closed = true;
1053 } else {
1054 discard_transaction_state(&mut inner);
1055 }
1056 }
1057 }
1058 }
1059}