reifydb_transaction/single/
mod.rs1use 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 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 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 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)); 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 barrier1.wait();
504
505 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 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 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 barrier1.wait();
557
558 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 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 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 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 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 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 });
729
730 handle1.join().unwrap();
731
732 thread::sleep(Duration::from_milliseconds(10).unwrap().to_std());
733
734 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}