reifydb_transaction/single/
write.rs1use 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}