Skip to main content

reifydb_transaction/single/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::Arc;
5
6use crossbeam_skiplist::SkipMap;
7use reifydb_codec::{
8	key::encoded::{EncodedKey, EncodedKeyRange},
9	row::bytes::EncodedBytes,
10};
11use reifydb_core::{delta::Delta, event::EventBus, interface::WithEventBus};
12use reifydb_runtime::sync::rwlock::{ArcRwLock, RwLock};
13use reifydb_store_single::SingleStore;
14
15pub mod read;
16pub mod write;
17
18use read::{KeyReadLock, SingleReadTransaction};
19use reifydb_runtime::{actor::system::ActorSystem, context::clock::Clock};
20use reifydb_value::Result;
21use write::{KeyWriteLock, SingleWriteTransaction};
22
23#[derive(Clone)]
24pub struct SingleTransaction {
25	inner: Arc<SingleTransactionInner>,
26}
27
28pub(crate) struct SingleTransactionInner {
29	pub(crate) store: RwLock<SingleStore>,
30	pub(crate) event_bus: EventBus,
31	pub(crate) key_locks: SkipMap<EncodedKey, ArcRwLock<()>>,
32}
33
34impl SingleTransactionInner {
35	fn get_or_create_lock(&self, key: &EncodedKey) -> ArcRwLock<()> {
36		if let Some(entry) = self.key_locks.get(key) {
37			return entry.value().clone();
38		}
39
40		self.key_locks.get_or_insert(key.clone(), ArcRwLock::new(())).value().clone()
41	}
42}
43
44impl SingleTransaction {
45	pub fn new(store: SingleStore, event_bus: EventBus) -> Self {
46		Self {
47			inner: Arc::new(SingleTransactionInner {
48				store: RwLock::new(store),
49				event_bus,
50				key_locks: SkipMap::new(),
51			}),
52		}
53	}
54
55	pub fn read_store(&self) -> SingleStore {
56		self.inner.store.read().clone()
57	}
58
59	pub fn testing() -> Self {
60		let actor_system = ActorSystem::testing(Clock::Real);
61		let spawner = actor_system.spawner();
62		Self::new(SingleStore::testing_memory(), EventBus::new(&spawner))
63	}
64
65	pub fn with_query<'a, I, F, R>(&self, keys: I, f: F) -> Result<R>
66	where
67		I: IntoIterator<Item = &'a EncodedKey>,
68		F: FnOnce(&mut SingleReadTransaction<'_>) -> Result<R>,
69	{
70		let mut tx = self.begin_query(keys)?;
71		f(&mut tx)
72	}
73
74	pub fn with_command<'a, I, F, R>(&self, keys: I, f: F) -> Result<R>
75	where
76		I: IntoIterator<Item = &'a EncodedKey>,
77		F: FnOnce(&mut SingleWriteTransaction<'_>) -> Result<R>,
78	{
79		let mut tx = self.begin_command(keys)?;
80		let result = f(&mut tx)?;
81		tx.commit()?;
82		Ok(result)
83	}
84
85	pub fn begin_query<'a, I>(&self, keys: I) -> Result<SingleReadTransaction<'_>>
86	where
87		I: IntoIterator<Item = &'a EncodedKey>,
88	{
89		let mut keys_vec: Vec<EncodedKey> = keys.into_iter().cloned().collect();
90		assert!(
91			!keys_vec.is_empty(),
92			"SVL transactions must declare keys upfront - empty keysets are not allowed"
93		);
94
95		keys_vec.sort();
96
97		let mut locks = Vec::new();
98		for key in &keys_vec {
99			let lock = self.inner.get_or_create_lock(key);
100			locks.push(KeyReadLock::new(lock));
101		}
102
103		Ok(SingleReadTransaction {
104			inner: &self.inner,
105			keys: keys_vec,
106			_key_locks: locks,
107		})
108	}
109
110	pub fn begin_command<'a, I>(&self, keys: I) -> Result<SingleWriteTransaction<'_>>
111	where
112		I: IntoIterator<Item = &'a EncodedKey>,
113	{
114		self.begin_command_ranged(keys, Vec::new())
115	}
116
117	pub fn begin_command_ranged<'a, I>(
118		&self,
119		lock_keys: I,
120		ranges: Vec<EncodedKeyRange>,
121	) -> Result<SingleWriteTransaction<'_>>
122	where
123		I: IntoIterator<Item = &'a EncodedKey>,
124	{
125		let mut keys_vec: Vec<EncodedKey> = lock_keys.into_iter().cloned().collect();
126		assert!(
127			!keys_vec.is_empty(),
128			"SVL transactions must declare keys upfront - empty keysets are not allowed"
129		);
130
131		keys_vec.sort();
132
133		let mut locks = Vec::new();
134		for key in &keys_vec {
135			let lock = self.inner.get_or_create_lock(key);
136			locks.push(KeyWriteLock::new(lock));
137		}
138
139		Ok(SingleWriteTransaction::new(&self.inner, keys_vec, ranges, locks))
140	}
141}
142
143impl WithEventBus for SingleTransaction {
144	fn event_bus(&self) -> &EventBus {
145		&self.inner.event_bus
146	}
147}
148
149#[cfg(test)]
150pub mod tests {
151	use std::{
152		iter,
153		ops::Bound,
154		sync::{Arc, Barrier},
155		thread,
156	};
157
158	use reifydb_core::{interface::catalog::id::QueueId, key::queue::QueueDeduplicationKey};
159	use reifydb_value::{util::cowvec::CowVec, value::duration::Duration};
160
161	use super::*;
162
163	fn make_key(s: &str) -> QueueDeduplicationKey {
164		QueueDeduplicationKey::new(QueueId(1), s.as_bytes().iter().map(|b| !b).collect::<Vec<u8>>())
165	}
166
167	fn make_value(s: &str) -> EncodedBytes {
168		EncodedBytes(CowVec::new(s.as_bytes().to_vec()))
169	}
170
171	fn create_test_svl() -> SingleTransaction {
172		SingleTransaction::testing()
173	}
174
175	#[test]
176	fn test_allowed_key_query() {
177		let svl = create_test_svl();
178		let key = make_key("test_key");
179
180		let mut tx = svl.begin_query(vec![&key.encode()]).unwrap();
181
182		let result = tx.get(&key.encode());
183		assert!(result.is_ok());
184	}
185
186	#[test]
187	fn test_disallowed_key_query() {
188		let svl = create_test_svl();
189		let key1 = make_key("allowed");
190		let key2 = make_key("disallowed");
191
192		let mut tx = svl.begin_query(vec![&key1.encode()]).unwrap();
193
194		assert!(tx.get(&key1.encode()).is_ok());
195
196		let result = tx.get(&key2.encode());
197		assert!(result.is_err());
198		let err = result.unwrap_err();
199		assert_eq!(err.0.code, "TXN_010");
200	}
201
202	#[test]
203	#[should_panic(expected = "SVL transactions must declare keys upfront - empty keysets are not allowed")]
204	fn test_empty_keyset_query_panics() {
205		let svl = create_test_svl();
206
207		let _tx = svl.begin_query(iter::empty());
208	}
209
210	#[test]
211	#[should_panic(expected = "SVL transactions must declare keys upfront - empty keysets are not allowed")]
212	fn test_empty_keyset_command_panics() {
213		let svl = create_test_svl();
214
215		let _tx = svl.begin_command(iter::empty());
216	}
217
218	#[test]
219	fn test_allowed_key_command() {
220		let svl = create_test_svl();
221		let key = make_key("test_key");
222		let value = make_value("test_value");
223
224		let mut tx = svl.begin_command(vec![&key.encode()]).unwrap();
225
226		assert!(tx.set(&key, value.clone()).is_ok());
227		assert!(tx.get(&key).is_ok());
228		assert!(tx.commit().is_ok());
229	}
230
231	#[test]
232	fn test_disallowed_key_command() {
233		let svl = create_test_svl();
234		let key1 = make_key("allowed");
235		let key2 = make_key("disallowed");
236		let value = make_value("test_value");
237
238		let mut tx = svl.begin_command(vec![&key1.encode()]).unwrap();
239
240		assert!(tx.set(&key1, value.clone()).is_ok());
241
242		let result = tx.set(&key2, value);
243		assert!(result.is_err());
244		let err = result.unwrap_err();
245		assert_eq!(err.0.code, "TXN_010");
246	}
247
248	#[test]
249	fn test_ranged_command_allows_keys_in_declared_range() {
250		let svl = create_test_svl();
251		let lock_key = make_key("lock");
252		let in_range = make_key("range_b");
253		let value = make_value("test_value");
254
255		// Only the coarse lock key is locked; writes are scoped to the declared range.
256		let range = EncodedKeyRange::new(
257			Bound::Included(make_key("range_a").encode()),
258			Bound::Excluded(make_key("range_z").encode()),
259		);
260		let mut tx = svl.begin_command_ranged(vec![&lock_key.encode()], vec![range]).unwrap();
261
262		assert!(tx.set(&in_range, value.clone()).is_ok());
263		assert!(tx.commit().is_ok());
264
265		let mut rx = svl.begin_query(vec![&in_range.encode()]).unwrap();
266		let row = rx.get(&in_range.encode()).unwrap().unwrap();
267		assert_eq!(row.bytes, value);
268	}
269
270	#[test]
271	fn test_ranged_command_rejects_keys_outside_range_and_lock_set() {
272		let svl = create_test_svl();
273		let lock_key = make_key("lock");
274		let outside = make_key("zzz_outside");
275		let value = make_value("test_value");
276
277		let range = EncodedKeyRange::new(
278			Bound::Included(make_key("range_a").encode()),
279			Bound::Excluded(make_key("range_z").encode()),
280		);
281		let mut tx = svl.begin_command_ranged(vec![&lock_key.encode()], vec![range]).unwrap();
282
283		let result = tx.set(&outside, value);
284		assert!(result.is_err());
285		let err = result.unwrap_err();
286		assert_eq!(err.0.code, "TXN_010");
287	}
288
289	#[test]
290	fn test_ranged_command_still_allows_exact_lock_keys() {
291		let svl = create_test_svl();
292		let lock_key = make_key("lock");
293		let value = make_value("test_value");
294
295		let range = EncodedKeyRange::new(
296			Bound::Included(make_key("range_a").encode()),
297			Bound::Excluded(make_key("range_z").encode()),
298		);
299		let mut tx = svl.begin_command_ranged(vec![&lock_key.encode()], vec![range]).unwrap();
300
301		// A declared lock key stays writable even though it falls outside the range.
302		assert!(tx.set(&lock_key, value).is_ok());
303		assert!(tx.commit().is_ok());
304	}
305
306	#[test]
307	fn test_command_commit_with_valid_keys() {
308		let svl = create_test_svl();
309		let key1 = make_key("key1");
310		let key2 = make_key("key2");
311		let value1 = make_value("value1");
312		let value2 = make_value("value2");
313
314		{
315			let mut tx = svl.begin_command(vec![&key1.encode(), &key2.encode()]).unwrap();
316			tx.set(&key1, value1.clone()).unwrap();
317			tx.set(&key2, value2.clone()).unwrap();
318			tx.commit().unwrap();
319		}
320
321		{
322			let mut tx = svl.begin_query(vec![&key1.encode(), &key2.encode()]).unwrap();
323			let result1 = tx.get(&key1.encode()).unwrap();
324			let result2 = tx.get(&key2.encode()).unwrap();
325			assert!(result1.is_some());
326			assert!(result2.is_some());
327			assert_eq!(result1.unwrap().bytes, value1);
328			assert_eq!(result2.unwrap().bytes, value2);
329		}
330	}
331
332	#[test]
333	fn test_rollback_with_scoped_keys() {
334		let svl = create_test_svl();
335		let key = make_key("test_key");
336		let value = make_value("test_value");
337
338		{
339			let mut tx = svl.begin_command(vec![&key.encode()]).unwrap();
340			tx.set(&key, value).unwrap();
341			tx.rollback().unwrap();
342		}
343
344		{
345			let mut tx = svl.begin_query(vec![&key.encode()]).unwrap();
346			let result = tx.get(&key.encode()).unwrap();
347			assert!(result.is_none());
348		}
349	}
350
351	#[test]
352	fn test_concurrent_reads() {
353		let svl = Arc::new(create_test_svl());
354		let key = make_key("shared_key");
355		let value = make_value("shared_value");
356
357		{
358			let mut tx = svl.begin_command(vec![&key.encode()]).unwrap();
359			tx.set(&key, value.clone()).unwrap();
360			tx.commit().unwrap();
361		}
362
363		let mut handles = vec![];
364		for _ in 0..5 {
365			let svl_clone = Arc::clone(&svl);
366			let key_clone = key.clone();
367			let value_clone = value.clone();
368
369			let handle = thread::spawn(move || {
370				let mut tx = svl_clone.begin_query(vec![&key_clone.encode()]).unwrap();
371				let result = tx.get(&key_clone.encode()).unwrap();
372				assert!(result.is_some());
373				assert_eq!(result.unwrap().bytes, value_clone);
374			});
375			handles.push(handle);
376		}
377
378		for handle in handles {
379			handle.join().unwrap();
380		}
381	}
382
383	#[test]
384	fn test_concurrent_writers_disjoint_keys() {
385		let svl = Arc::new(create_test_svl());
386
387		let mut handles = vec![];
388		for i in 0..5 {
389			let svl_clone = Arc::clone(&svl);
390			let key = make_key(&format!("key_{}", i));
391			let value = make_value(&format!("value_{}", i));
392
393			let handle = thread::spawn(move || {
394				let mut tx = svl_clone.begin_command(vec![&key.encode()]).unwrap();
395				tx.set(&key, value).unwrap();
396				tx.commit().unwrap();
397			});
398			handles.push(handle);
399		}
400
401		for handle in handles {
402			handle.join().unwrap();
403		}
404
405		for i in 0..5 {
406			let key = make_key(&format!("key_{}", i));
407			let expected_value = make_value(&format!("value_{}", i));
408
409			let mut tx = svl.begin_query(vec![&key.encode()]).unwrap();
410			let result = tx.get(&key.encode()).unwrap();
411			assert!(result.is_some());
412			assert_eq!(result.unwrap().bytes, expected_value);
413		}
414	}
415
416	#[test]
417	fn test_concurrent_readers_and_writer() {
418		let svl = Arc::new(create_test_svl());
419		let key1 = make_key("key1");
420		let key2 = make_key("key2");
421		let value1 = make_value("value1");
422		let value2 = make_value("value2");
423
424		{
425			let mut tx = svl.begin_command(vec![&key1.encode(), &key2.encode()]).unwrap();
426			tx.set(&key1, value1.clone()).unwrap();
427			tx.set(&key2, value2.clone()).unwrap();
428			tx.commit().unwrap();
429		}
430
431		let mut handles = vec![];
432		for _ in 0..3 {
433			let svl_clone = Arc::clone(&svl);
434			let key_clone = key1.clone();
435			let value_clone = value1.clone();
436
437			let handle = thread::spawn(move || {
438				let mut tx = svl_clone.begin_query(vec![&key_clone.encode()]).unwrap();
439				let result = tx.get(&key_clone.encode()).unwrap();
440				assert!(result.is_some());
441				assert_eq!(result.unwrap().bytes, value_clone);
442			});
443			handles.push(handle);
444		}
445
446		// A writer on a different key must not block these readers.
447		let svl_clone = Arc::clone(&svl);
448		let new_value = make_value("new_value2");
449		let handle = thread::spawn(move || {
450			let mut tx = svl_clone.begin_command(vec![&key2.encode()]).unwrap();
451			tx.set(&key2, new_value).unwrap();
452			tx.commit().unwrap();
453		});
454		handles.push(handle);
455
456		for handle in handles {
457			handle.join().unwrap();
458		}
459	}
460
461	#[test]
462	fn test_no_panics_with_rwlock() {
463		let svl = Arc::new(create_test_svl());
464
465		let mut handles = vec![];
466		for i in 0..10 {
467			let svl_clone = Arc::clone(&svl);
468			let key = make_key(&format!("key_{}", i % 3)); // keys overlap across threads
469			let value = make_value(&format!("value_{}", i));
470
471			let handle = thread::spawn(move || {
472				if i % 2 == 0 {
473					let mut tx = svl_clone.begin_command(vec![&key.encode()]).unwrap();
474					let _ = tx.set(&key, value);
475					let _ = tx.commit();
476				} else {
477					let mut tx = svl_clone.begin_query(vec![&key.encode()]).unwrap();
478					let _ = tx.get(&key.encode());
479				}
480			});
481			handles.push(handle);
482		}
483
484		for handle in handles {
485			handle.join().unwrap();
486		}
487	}
488
489	#[test]
490	fn test_write_blocks_concurrent_write() {
491		let svl = Arc::new(create_test_svl());
492		let key = make_key("blocking_key");
493		let barrier = Arc::new(Barrier::new(2));
494
495		let svl1 = Arc::clone(&svl);
496		let key1 = key.clone();
497		let barrier1 = Arc::clone(&barrier);
498		let handle1 = thread::spawn(move || {
499			let mut tx = svl1.begin_command(vec![&key1.encode()]).unwrap();
500			tx.set(&key1, make_value("value1")).unwrap();
501
502			// Signal that the write lock is held.
503			barrier1.wait();
504
505			// Hold the lock so the second writer has to block on it.
506			thread::sleep(Duration::from_milliseconds(100).unwrap().to_std());
507
508			tx.commit().unwrap();
509		});
510
511		let svl2 = Arc::clone(&svl);
512		let key2 = key.clone();
513		let barrier2 = Arc::clone(&barrier);
514		let handle2 = thread::spawn(move || {
515			barrier2.wait();
516
517			// Give thread 1 time to enter its sleep still holding the lock.
518			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
519
520			let mut tx = svl2.begin_command(vec![&key2.encode()]).unwrap();
521			tx.set(&key2, make_value("value2")).unwrap();
522			tx.commit().unwrap();
523		});
524
525		handle1.join().unwrap();
526		handle2.join().unwrap();
527
528		// Thread 2 could only have written after thread 1 released, so its value must survive.
529		let mut tx = svl.begin_query(vec![&key.encode()]).unwrap();
530		let result = tx.get(&key.encode()).unwrap();
531		assert!(result.is_some());
532		assert_eq!(result.unwrap().bytes, make_value("value2"));
533	}
534
535	#[test]
536	fn test_write_blocks_concurrent_read() {
537		let svl = Arc::new(create_test_svl());
538		let key = make_key("blocking_key");
539
540		{
541			let mut tx = svl.begin_command(vec![&key.encode()]).unwrap();
542			tx.set(&key, make_value("initial")).unwrap();
543			tx.commit().unwrap();
544		}
545
546		let barrier = Arc::new(Barrier::new(2));
547
548		let svl1 = Arc::clone(&svl);
549		let key1 = key.clone();
550		let barrier1 = Arc::clone(&barrier);
551		let handle1 = thread::spawn(move || {
552			let mut tx = svl1.begin_command(vec![&key1.encode()]).unwrap();
553			tx.set(&key1, make_value("updated")).unwrap();
554
555			// Signal that the write lock is held.
556			barrier1.wait();
557
558			// Hold the lock so the reader has to block on it.
559			thread::sleep(Duration::from_milliseconds(100).unwrap().to_std());
560
561			tx.commit().unwrap();
562		});
563
564		let svl2 = Arc::clone(&svl);
565		let key2 = key.clone();
566		let barrier2 = Arc::clone(&barrier);
567		let handle2 = thread::spawn(move || {
568			barrier2.wait();
569
570			// Give thread 1 time to enter its sleep still holding the lock.
571			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
572
573			let mut tx = svl2.begin_query(vec![&key2.encode()]).unwrap();
574			let result = tx.get(&key2.encode()).unwrap();
575
576			// A reader that truly blocked sees the committed value, never the initial one.
577			assert!(result.is_some());
578			assert_eq!(result.unwrap().bytes, make_value("updated"));
579		});
580
581		handle1.join().unwrap();
582		handle2.join().unwrap();
583	}
584
585	#[test]
586	fn test_concurrent_reads_allowed() {
587		let svl = Arc::new(create_test_svl());
588		let key = make_key("shared_read_key");
589
590		{
591			let mut tx = svl.begin_command(vec![&key.encode()]).unwrap();
592			tx.set(&key, make_value("shared")).unwrap();
593			tx.commit().unwrap();
594		}
595
596		let barrier = Arc::new(Barrier::new(3));
597		let mut handles = vec![];
598
599		for _ in 0..3 {
600			let svl_clone = Arc::clone(&svl);
601			let key_clone = key.clone();
602			let barrier_clone = Arc::clone(&barrier);
603
604			let handle = thread::spawn(move || {
605				let mut tx = svl_clone.begin_query(vec![&key_clone.encode()]).unwrap();
606
607				// Every reader holds its read lock past this point; a mutually exclusive
608				// read lock would deadlock here rather than fail an assertion.
609				barrier_clone.wait();
610
611				let result = tx.get(&key_clone.encode()).unwrap();
612				assert!(result.is_some());
613				assert_eq!(result.unwrap().bytes, make_value("shared"));
614
615				thread::sleep(Duration::from_milliseconds(50).unwrap().to_std());
616			});
617			handles.push(handle);
618		}
619
620		for handle in handles {
621			handle.join().unwrap();
622		}
623	}
624
625	#[test]
626	fn test_overlapping_keys_different_order() {
627		// The two threads declare the same keys in opposite order; begin_command sorts them, so
628		// there is no lock-order cycle to deadlock on.
629		let svl = Arc::new(create_test_svl());
630		let key1 = make_key("deadlock_key1");
631		let key2 = make_key("deadlock_key2");
632		let barrier = Arc::new(Barrier::new(2));
633
634		let svl1 = Arc::clone(&svl);
635		let key1_clone = key1.clone();
636		let key2_clone = key2.clone();
637		let barrier1 = Arc::clone(&barrier);
638		let handle1 = thread::spawn(move || {
639			barrier1.wait();
640			let mut tx = svl1.begin_command(vec![&key1_clone.encode(), &key2_clone.encode()]).unwrap();
641			tx.set(&key1_clone, make_value("from_thread1")).unwrap();
642			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
643			tx.commit().unwrap();
644		});
645
646		let svl2 = Arc::clone(&svl);
647		let key1_clone2 = key1.clone();
648		let key2_clone2 = key2.clone();
649		let barrier2 = Arc::clone(&barrier);
650		let handle2 = thread::spawn(move || {
651			barrier2.wait();
652			let mut tx = svl2.begin_command(vec![&key2_clone2.encode(), &key1_clone2.encode()]).unwrap();
653			tx.set(&key2_clone2, make_value("from_thread2")).unwrap();
654			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
655			tx.commit().unwrap();
656		});
657
658		handle1.join().unwrap();
659		handle2.join().unwrap();
660
661		let mut tx = svl.begin_query(vec![&key1.encode(), &key2.encode()]).unwrap();
662		let result1 = tx.get(&key1.encode()).unwrap();
663		let result2 = tx.get(&key2.encode()).unwrap();
664		assert!(result1.is_some());
665		assert!(result2.is_some());
666	}
667
668	#[test]
669	fn test_circular_dependency_three_transactions() {
670		// The three key sets chain into a cycle (1,2) (2,3) (3,1); only sorted acquisition keeps
671		// that from becoming a deadlock.
672		let svl = Arc::new(create_test_svl());
673		let key1 = make_key("circular_key1");
674		let key2 = make_key("circular_key2");
675		let key3 = make_key("circular_key3");
676		let barrier = Arc::new(Barrier::new(3));
677
678		let svl1 = Arc::clone(&svl);
679		let k1_1 = key1.clone();
680		let k2_1 = key2.clone();
681		let barrier1 = Arc::clone(&barrier);
682		let handle1 = thread::spawn(move || {
683			barrier1.wait();
684			let mut tx = svl1.begin_command(vec![&k1_1.encode(), &k2_1.encode()]).unwrap();
685			tx.set(&k1_1, make_value("t1")).unwrap();
686			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
687			tx.commit().unwrap();
688		});
689
690		let svl2 = Arc::clone(&svl);
691		let k2_2 = key2.clone();
692		let k3_2 = key3.clone();
693		let barrier2 = Arc::clone(&barrier);
694		let handle2 = thread::spawn(move || {
695			barrier2.wait();
696			let mut tx = svl2.begin_command(vec![&k2_2.encode(), &k3_2.encode()]).unwrap();
697			tx.set(&k2_2, make_value("t2")).unwrap();
698			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
699			tx.commit().unwrap();
700		});
701
702		let svl3 = Arc::clone(&svl);
703		let barrier3 = Arc::clone(&barrier);
704		let handle3 = thread::spawn(move || {
705			barrier3.wait();
706			let mut tx = svl3.begin_command(vec![&key3.encode(), &key1.encode()]).unwrap();
707			tx.set(&key3, make_value("t3")).unwrap();
708			thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
709			tx.commit().unwrap();
710		});
711
712		handle1.join().unwrap();
713		handle2.join().unwrap();
714		handle3.join().unwrap();
715	}
716
717	#[test]
718	fn test_locks_released_on_drop() {
719		let svl = Arc::new(create_test_svl());
720		let key = make_key("drop_test_key");
721
722		let svl1 = Arc::clone(&svl);
723		let key_clone = key.clone();
724		let handle1 = thread::spawn(move || {
725			let mut tx = svl1.begin_command(vec![&key_clone.encode()]).unwrap();
726			tx.set(&key_clone, make_value("dropped")).unwrap();
727			// Dropped here without commit.
728		});
729
730		handle1.join().unwrap();
731
732		thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
733
734		// A lock not released on drop makes this block forever rather than fail.
735		let svl2 = Arc::clone(&svl);
736		let key_clone2 = key.clone();
737		let handle2 = thread::spawn(move || {
738			let mut tx = svl2.begin_command(vec![&key_clone2.encode()]).unwrap();
739			tx.set(&key_clone2, make_value("success")).unwrap();
740			tx.commit().unwrap();
741		});
742
743		handle2.join().unwrap();
744
745		let mut tx = svl.begin_query(vec![&key.encode()]).unwrap();
746		let result = tx.get(&key.encode()).unwrap();
747		assert!(result.is_some());
748		assert_eq!(result.unwrap().bytes, make_value("success"));
749	}
750}