Skip to main content

reifydb_transaction/single/
write.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{mem::take, ops::RangeBounds};
5
6use indexmap::IndexMap;
7use reifydb_core::{
8	interface::store::{SingleVersionCommit, SingleVersionContains, SingleVersionGet, SingleVersionRow},
9	key::any::TaggedKey,
10};
11use reifydb_runtime::sync::rwlock::{ArcRwLock, OwnedRwLockWriteGuard};
12use reifydb_value::{
13	Result, reifydb_assertions,
14	util::{cowvec::CowVec, hex::encode},
15};
16
17use super::*;
18use crate::error::TransactionError;
19
20pub struct KeyWriteLock {
21	pub(super) _guard: OwnedRwLockWriteGuard<()>,
22}
23
24impl KeyWriteLock {
25	pub(super) fn new(lock: ArcRwLock<()>) -> Self {
26		Self {
27			_guard: lock.write(),
28		}
29	}
30}
31
32pub struct SingleWriteTransaction<'a> {
33	pub(super) inner: &'a SingleTransactionInner,
34	pub(super) keys: Vec<EncodedKey>,
35	pub(super) ranges: Vec<EncodedKeyRange>,
36	pub(super) _key_locks: Vec<KeyWriteLock>,
37	pub(super) pending: IndexMap<EncodedKey, Delta>,
38	pub(super) completed: bool,
39}
40
41impl<'a> SingleWriteTransaction<'a> {
42	pub(super) fn new(
43		inner: &'a SingleTransactionInner,
44		keys: Vec<EncodedKey>,
45		ranges: Vec<EncodedKeyRange>,
46		key_locks: Vec<KeyWriteLock>,
47	) -> Self {
48		Self {
49			inner,
50			keys,
51			ranges,
52			_key_locks: key_locks,
53			pending: IndexMap::new(),
54			completed: false,
55		}
56	}
57
58	#[inline]
59	fn check_key_allowed(&self, key: &EncodedKey) -> Result<()> {
60		if self.keys.iter().any(|k| k == key) || self.ranges.iter().any(|range| range.contains(key)) {
61			Ok(())
62		} else {
63			Err(TransactionError::KeyOutOfScope {
64				key: encode(key),
65			}
66			.into())
67		}
68	}
69
70	pub fn get<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<Option<SingleVersionRow>> {
71		let encoded = key.clone().into().encode();
72		self.check_key_allowed(&encoded)?;
73
74		if let Some(delta) = self.pending.get(&encoded) {
75			return match delta {
76				Delta::Set {
77					bytes,
78					..
79				} => Ok(Some(SingleVersionRow {
80					key: encoded,
81					bytes: bytes.clone(),
82				})),
83				Delta::Remove {
84					..
85				} => Ok(None),
86			};
87		}
88
89		let store = self.inner.store.read().clone();
90		SingleVersionGet::get(&store, &encoded)
91	}
92
93	pub fn contains_key<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<bool> {
94		let key = &key.clone().into().encode();
95		self.check_key_allowed(key)?;
96
97		if let Some(delta) = self.pending.get(key) {
98			return match delta {
99				Delta::Set {
100					..
101				} => Ok(true),
102				Delta::Remove {
103					..
104				} => Ok(false),
105			};
106		}
107
108		let store = self.inner.store.read().clone();
109		SingleVersionContains::contains(&store, key)
110	}
111
112	pub fn set<K: Into<TaggedKey> + Clone>(&mut self, key: &K, bytes: impl Into<EncodedBytes>) -> Result<()> {
113		let key: TaggedKey = key.clone().into();
114		let encoded = key.encode();
115		self.check_key_allowed(&encoded)?;
116
117		let delta = Delta::Set {
118			key,
119			bytes: bytes.into(),
120		};
121		self.pending.insert(encoded, delta);
122		Ok(())
123	}
124
125	pub fn remove_with_pre<K: Into<TaggedKey> + Clone>(&mut self, key: &K, pre: EncodedBytes) -> Result<()> {
126		let key: TaggedKey = key.clone().into();
127		let encoded = key.encode();
128		self.check_key_allowed(&encoded)?;
129
130		self.pending.insert(encoded, Delta::remove_announced(key, pre));
131		Ok(())
132	}
133
134	pub fn remove<K: Into<TaggedKey> + Clone>(&mut self, key: &K) -> Result<()> {
135		let key: TaggedKey = key.clone().into();
136		let encoded = key.encode();
137		self.check_key_allowed(&encoded)?;
138
139		self.pending.insert(encoded, Delta::remove_silent(key));
140		Ok(())
141	}
142
143	pub fn commit(&mut self) -> Result<()> {
144		let deltas = self.drain_pending();
145
146		if !deltas.is_empty() {
147			self.commit_deltas(deltas)?;
148		}
149
150		self.completed = true;
151		Ok(())
152	}
153
154	#[inline]
155	fn drain_pending(&mut self) -> Vec<Delta> {
156		take(&mut self.pending).into_iter().map(|(_, delta)| delta).collect()
157	}
158
159	#[inline]
160	fn commit_deltas(&self, deltas: Vec<Delta>) -> Result<()> {
161		reifydb_assertions! {
162			let count = deltas.len();
163			assert!(
164				count > 0,
165				"commit_deltas must not run on an empty delta set; an empty store commit \
166				 acquires the store write lock for a no-op transaction (count={count})"
167			);
168		}
169
170		let mut store = self.inner.store.write();
171		SingleVersionCommit::commit(&mut *store, CowVec::new(deltas))
172	}
173
174	pub fn rollback(&mut self) -> Result<()> {
175		self.pending.clear();
176		self.completed = true;
177		Ok(())
178	}
179}