Skip to main content

reifydb_sub_flow/transaction/
write.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_codec::{encoded::row::EncodedRow, key::encoded::EncodedKey};
5use reifydb_value::Result;
6
7use super::FlowTransaction;
8
9impl FlowTransaction {
10	pub fn set(&mut self, key: &EncodedKey, value: EncodedRow) -> Result<()> {
11		match self {
12			Self::Committing {
13				cmd,
14				..
15			} => cmd.set(key, value),
16			_ => {
17				self.inner_mut().pending.insert(key.clone(), value);
18				Ok(())
19			}
20		}
21	}
22
23	pub fn remove(&mut self, key: &EncodedKey) -> Result<()> {
24		match self {
25			Self::Committing {
26				cmd,
27				..
28			} => cmd.remove(key),
29			_ => {
30				self.inner_mut().pending.remove(key.clone());
31				Ok(())
32			}
33		}
34	}
35
36	pub fn drop_key(&mut self, key: &EncodedKey) -> Result<()> {
37		match self {
38			Self::Committing {
39				cmd,
40				..
41			} => cmd.drop_key(key),
42			_ => {
43				self.inner_mut().pending.drop_key(key.clone());
44				Ok(())
45			}
46		}
47	}
48
49	pub fn set_batch(&mut self, keys: &[EncodedKey], values: &[EncodedRow]) -> Result<()> {
50		match self {
51			Self::Committing {
52				cmd,
53				..
54			} => {
55				for (key, value) in keys.iter().zip(values.iter()) {
56					cmd.set(key, value.clone())?;
57				}
58				Ok(())
59			}
60			_ => {
61				self.inner_mut().pending.insert_batch(keys, values);
62				Ok(())
63			}
64		}
65	}
66
67	pub fn remove_batch(&mut self, keys: &[EncodedKey]) -> Result<()> {
68		match self {
69			Self::Committing {
70				cmd,
71				..
72			} => {
73				for key in keys {
74					cmd.remove(key)?;
75				}
76				Ok(())
77			}
78			_ => {
79				self.inner_mut().pending.remove_batch(keys);
80				Ok(())
81			}
82		}
83	}
84
85	pub fn drop_keys(&mut self, keys: &[EncodedKey]) -> Result<()> {
86		match self {
87			Self::Committing {
88				cmd,
89				..
90			} => {
91				for key in keys {
92					cmd.drop_key(key)?;
93				}
94				Ok(())
95			}
96			_ => {
97				self.inner_mut().pending.drop_keys(keys);
98				Ok(())
99			}
100		}
101	}
102}
103
104#[cfg(test)]
105pub mod tests {
106	use reifydb_catalog::catalog::Catalog;
107	use reifydb_codec::{encoded::row::EncodedRow, key::encoded::EncodedKey};
108	use reifydb_core::common::CommitVersion;
109	use reifydb_runtime::context::clock::{Clock, MockClock};
110	use reifydb_transaction::{interceptor::interceptors::Interceptors, transaction::admin::AdminTransaction};
111	use reifydb_value::util::cowvec::CowVec;
112
113	use super::*;
114	use crate::operator::stateful::test_utils::test::create_test_transaction;
115
116	fn make_key(s: &str) -> EncodedKey {
117		EncodedKey::new(s.as_bytes().to_vec())
118	}
119
120	fn make_value(s: &str) -> EncodedRow {
121		EncodedRow(CowVec::new(s.as_bytes().to_vec()))
122	}
123
124	fn get_row(parent: &mut AdminTransaction, key: &EncodedKey) -> Option<EncodedRow> {
125		parent.get(key).unwrap().map(|m| m.row.clone())
126	}
127
128	#[test]
129	fn test_set_buffers_to_pending() {
130		let parent = create_test_transaction();
131		let mut txn = FlowTransaction::deferred(
132			&parent,
133			CommitVersion(1),
134			Catalog::testing(),
135			Interceptors::new(),
136			Clock::Mock(MockClock::from_millis(1000)),
137		);
138
139		let key = make_key("key1");
140		let value = make_value("value1");
141
142		txn.set(&key, value.clone()).unwrap();
143
144		// Value should be in pending buffer
145		assert_eq!(txn.pending().get(&key), Some(&value));
146	}
147
148	#[test]
149	fn test_set_multiple_keys() {
150		let parent = create_test_transaction();
151		let mut txn = FlowTransaction::deferred(
152			&parent,
153			CommitVersion(1),
154			Catalog::testing(),
155			Interceptors::new(),
156			Clock::Mock(MockClock::from_millis(1000)),
157		);
158
159		txn.set(&make_key("key1"), make_value("value1")).unwrap();
160		txn.set(&make_key("key2"), make_value("value2")).unwrap();
161		txn.set(&make_key("key3"), make_value("value3")).unwrap();
162
163		assert_eq!(txn.pending().get(&make_key("key1")), Some(&make_value("value1")));
164		assert_eq!(txn.pending().get(&make_key("key2")), Some(&make_value("value2")));
165		assert_eq!(txn.pending().get(&make_key("key3")), Some(&make_value("value3")));
166	}
167
168	#[test]
169	fn test_set_overwrites_same_key() {
170		let parent = create_test_transaction();
171		let mut txn = FlowTransaction::deferred(
172			&parent,
173			CommitVersion(1),
174			Catalog::testing(),
175			Interceptors::new(),
176			Clock::Mock(MockClock::from_millis(1000)),
177		);
178
179		let key = make_key("key1");
180		txn.set(&key, make_value("value1")).unwrap();
181		txn.set(&key, make_value("value2")).unwrap();
182
183		// Should have only one entry with latest value
184		assert_eq!(txn.pending().get(&key), Some(&make_value("value2")));
185	}
186
187	#[test]
188	fn test_remove_buffers_to_pending() {
189		let parent = create_test_transaction();
190		let mut txn = FlowTransaction::deferred(
191			&parent,
192			CommitVersion(1),
193			Catalog::testing(),
194			Interceptors::new(),
195			Clock::Mock(MockClock::from_millis(1000)),
196		);
197
198		let key = make_key("key1");
199		txn.remove(&key).unwrap();
200
201		// Key should be marked for removal in pending buffer
202		assert!(txn.pending().is_removed(&key));
203	}
204
205	#[test]
206	fn test_remove_multiple_keys() {
207		let parent = create_test_transaction();
208		let mut txn = FlowTransaction::deferred(
209			&parent,
210			CommitVersion(1),
211			Catalog::testing(),
212			Interceptors::new(),
213			Clock::Mock(MockClock::from_millis(1000)),
214		);
215
216		txn.remove(&make_key("key1")).unwrap();
217		txn.remove(&make_key("key2")).unwrap();
218		txn.remove(&make_key("key3")).unwrap();
219
220		assert!(txn.pending().is_removed(&make_key("key1")));
221		assert!(txn.pending().is_removed(&make_key("key2")));
222		assert!(txn.pending().is_removed(&make_key("key3")));
223	}
224
225	#[test]
226	fn test_set_then_remove() {
227		let parent = create_test_transaction();
228		let mut txn = FlowTransaction::deferred(
229			&parent,
230			CommitVersion(1),
231			Catalog::testing(),
232			Interceptors::new(),
233			Clock::Mock(MockClock::from_millis(1000)),
234		);
235
236		let key = make_key("key1");
237		txn.set(&key, make_value("value1")).unwrap();
238		assert_eq!(txn.pending().get(&key), Some(&make_value("value1")));
239
240		txn.remove(&key).unwrap();
241		assert!(txn.pending().is_removed(&key));
242		assert_eq!(txn.pending().get(&key), None);
243	}
244
245	#[test]
246	fn test_remove_then_set() {
247		let parent = create_test_transaction();
248		let mut txn = FlowTransaction::deferred(
249			&parent,
250			CommitVersion(1),
251			Catalog::testing(),
252			Interceptors::new(),
253			Clock::Mock(MockClock::from_millis(1000)),
254		);
255
256		let key = make_key("key1");
257		txn.remove(&key).unwrap();
258		assert!(txn.pending().is_removed(&key));
259
260		txn.set(&key, make_value("value1")).unwrap();
261		assert!(!txn.pending().is_removed(&key));
262		assert_eq!(txn.pending().get(&key), Some(&make_value("value1")));
263	}
264
265	#[test]
266	fn test_writes_not_visible_to_parent() {
267		let mut parent = create_test_transaction();
268		let mut txn = FlowTransaction::deferred(
269			&parent,
270			CommitVersion(1),
271			Catalog::testing(),
272			Interceptors::new(),
273			Clock::Mock(MockClock::from_millis(1000)),
274		);
275
276		let key = make_key("key1");
277		let value = make_value("value1");
278
279		// Set in FlowTransaction
280		txn.set(&key, value.clone()).unwrap();
281
282		// Parent should not see the write
283		assert_eq!(get_row(&mut parent, &key), None);
284	}
285
286	#[test]
287	fn test_removes_not_visible_to_parent() {
288		let mut parent = create_test_transaction();
289
290		// Set a value in parent
291		let key = make_key("key1");
292		let value = make_value("value1");
293		parent.set(&key, value.clone()).unwrap();
294		assert_eq!(get_row(&mut parent, &key), Some(value.clone()));
295
296		// Create FlowTransaction and remove the key
297		let parent_version = parent.version();
298		let mut txn = FlowTransaction::deferred(
299			&parent,
300			parent_version,
301			Catalog::testing(),
302			Interceptors::new(),
303			Clock::Mock(MockClock::from_millis(1000)),
304		);
305		txn.remove(&key).unwrap();
306
307		// Parent should still see the value
308		assert_eq!(get_row(&mut parent, &key), Some(value));
309	}
310
311	#[test]
312	fn test_mixed_writes_and_removes() {
313		let parent = create_test_transaction();
314		let mut txn = FlowTransaction::deferred(
315			&parent,
316			CommitVersion(1),
317			Catalog::testing(),
318			Interceptors::new(),
319			Clock::Mock(MockClock::from_millis(1000)),
320		);
321
322		txn.set(&make_key("write1"), make_value("v1")).unwrap();
323		txn.remove(&make_key("remove1")).unwrap();
324		txn.set(&make_key("write2"), make_value("v2")).unwrap();
325		txn.remove(&make_key("remove2")).unwrap();
326
327		assert_eq!(txn.pending().get(&make_key("write1")), Some(&make_value("v1")));
328		assert_eq!(txn.pending().get(&make_key("write2")), Some(&make_value("v2")));
329		assert!(txn.pending().is_removed(&make_key("remove1")));
330		assert!(txn.pending().is_removed(&make_key("remove2")));
331	}
332}