iota_core/execution_cache/
object_locks.rs1use dashmap::{DashMap, mapref::entry::Entry as DashMapEntry};
6use iota_common::debug_fatal;
7use iota_sdk_types::{ObjectId, ObjectReference};
8use iota_types::{
9 error::{IotaError, IotaResult, UserInputError},
10 object::Object,
11 storage::ObjectStore,
12 transaction::VerifiedSignedTransaction,
13};
14use tracing::{debug, info, instrument, trace};
15
16use super::writeback_cache::WritebackCache;
17use crate::authority::authority_per_epoch_store::{AuthorityPerEpochStore, LockDetails};
18
19type RefCount = usize;
20
21pub(super) struct ObjectLocks {
22 locked_transactions: DashMap<ObjectReference, (RefCount, LockDetails)>,
33}
34
35impl ObjectLocks {
36 pub fn new() -> Self {
37 Self {
38 locked_transactions: DashMap::new(),
39 }
40 }
41
42 pub(crate) fn get_transaction_lock(
43 &self,
44 obj_ref: &ObjectReference,
45 epoch_store: &AuthorityPerEpochStore,
46 ) -> IotaResult<Option<LockDetails>> {
47 epoch_store.tables()?.get_locked_transaction(obj_ref)
51 }
52
53 fn try_set_transaction_lock(
59 &self,
60 obj_ref: &ObjectReference,
61 new_lock: LockDetails,
62 epoch_store: &AuthorityPerEpochStore,
63 ) -> IotaResult {
64 let entry = self.locked_transactions.entry(*obj_ref);
66
67 let conflict = |prev_lock: LockDetails| {
85 debug!(
86 "lock conflict detected for {:?}: {:?} != {:?}",
87 obj_ref, prev_lock, new_lock
88 );
89 IotaError::ObjectLockConflict {
90 obj_ref: *obj_ref,
91 pending_transaction: prev_lock,
92 }
93 };
94
95 match entry {
99 DashMapEntry::Vacant(vacant) => {
100 let tables = epoch_store.tables()?;
101 match tables.get_locked_transaction(obj_ref)? {
102 Some(prev_lock) if prev_lock != new_lock => return Err(conflict(prev_lock)),
103 Some(prev_lock) => trace!("read lock from db: {:?}", prev_lock),
104 None => trace!("set lock: {:?}", new_lock),
105 }
106 vacant.insert((1, new_lock));
107 }
108 DashMapEntry::Occupied(mut occupied) => {
109 let prev_lock = occupied.get().1;
110 if prev_lock != new_lock {
111 return Err(conflict(prev_lock));
112 }
113 occupied.get_mut().0 += 1;
114 }
115 }
116
117 Ok(())
118 }
119
120 pub(crate) fn clear(&self) {
121 info!("clearing old transaction locks");
122 self.locked_transactions.clear();
123 }
124
125 fn verify_live_object(obj_ref: &ObjectReference, live_object: &Object) -> IotaResult {
126 debug_assert_eq!(obj_ref.object_id, live_object.id());
127 if obj_ref.version != live_object.version() {
128 debug!(
129 "object version unavailable for consumption: {:?} (current: {})",
130 obj_ref,
131 live_object.version()
132 );
133 return Err(IotaError::UserInput {
134 error: UserInputError::ObjectVersionUnavailableForConsumption {
135 provided_obj_ref: *obj_ref,
136 current_version: live_object.version(),
137 },
138 });
139 }
140
141 let live_digest = live_object.digest();
142 if obj_ref.digest != live_digest {
143 return Err(IotaError::UserInput {
144 error: UserInputError::InvalidObjectDigest {
145 object_id: obj_ref.object_id,
146 expected_digest: live_digest,
147 },
148 });
149 }
150
151 Ok(())
152 }
153
154 fn clear_cached_locks(&self, locks: &[(ObjectReference, LockDetails)]) {
155 for (obj_ref, lock) in locks {
156 let entry = self.locked_transactions.entry(*obj_ref);
157 let mut occupied = match entry {
158 DashMapEntry::Vacant(_) => {
159 debug_fatal!("lock must exist for object: {:?}", obj_ref);
160 continue;
161 }
162 DashMapEntry::Occupied(occupied) => occupied,
163 };
164
165 if occupied.get().1 == *lock {
166 occupied.get_mut().0 -= 1;
167 if occupied.get().0 == 0 {
168 trace!("clearing lock: {:?}", lock);
169 occupied.remove();
170 }
171 } else {
172 panic!("lock was changed since we set it");
176 }
177 }
178 }
179
180 fn multi_get_objects_must_exist(
181 cache: &WritebackCache,
182 object_ids: &[ObjectId],
183 ) -> IotaResult<Vec<Object>> {
184 let objects = cache.try_multi_get_objects(object_ids)?;
185 let mut result = Vec::with_capacity(objects.len());
186 for (i, object) in objects.into_iter().enumerate() {
187 if let Some(object) = object {
188 result.push(object);
189 } else {
190 return Err(IotaError::UserInput {
191 error: UserInputError::ObjectNotFound {
192 object_id: object_ids[i],
193 version: None,
194 },
195 });
196 }
197 }
198 Ok(result)
199 }
200
201 pub(crate) fn validate_owned_object_versions(
204 cache: &WritebackCache,
205 owned_input_objects: &[ObjectReference],
206 ) -> IotaResult {
207 if owned_input_objects.is_empty() {
208 return Ok(());
209 }
210 let object_ids: Vec<_> = owned_input_objects.iter().map(|o| o.object_id).collect();
211 let live_objects = Self::multi_get_objects_must_exist(cache, &object_ids)?;
212 for (obj_ref, live_object) in owned_input_objects.iter().zip(live_objects.iter()) {
213 Self::verify_live_object(obj_ref, live_object)?;
214 }
215 Ok(())
216 }
217
218 #[instrument(level = "debug", skip_all)]
219 pub(crate) fn acquire_transaction_locks(
220 &self,
221 cache: &WritebackCache,
222 epoch_store: &AuthorityPerEpochStore,
223 owned_input_objects: &[ObjectReference],
224 transaction: VerifiedSignedTransaction,
225 ) -> IotaResult {
226 let tx_digest = *transaction.digest();
227
228 let object_ids = owned_input_objects
229 .iter()
230 .map(|o| o.object_id)
231 .collect::<Vec<_>>();
232 let live_objects = Self::multi_get_objects_must_exist(cache, &object_ids)?;
233
234 for (obj_ref, live_object) in owned_input_objects.iter().zip(live_objects.iter()) {
236 Self::verify_live_object(obj_ref, live_object)?;
237 }
238
239 let mut locks_to_write: Vec<(_, LockDetails)> =
240 Vec::with_capacity(owned_input_objects.len());
241
242 let owned_input_objects = {
253 let mut o = owned_input_objects.to_vec();
254 o.sort_by_key(|o| o.object_id);
255 o
256 };
257
258 for obj_ref in owned_input_objects.iter() {
265 match self.try_set_transaction_lock(obj_ref, tx_digest, epoch_store) {
266 Ok(()) => locks_to_write.push((*obj_ref, tx_digest)),
267 Err(e) => {
268 self.clear_cached_locks(&locks_to_write);
274 return Err(e);
275 }
276 }
277 }
278
279 epoch_store
281 .tables()?
282 .write_transaction_locks(transaction, locks_to_write.iter().cloned())?;
283
284 self.clear_cached_locks(&locks_to_write);
286
287 Ok(())
288 }
289}
290
291#[cfg(test)]
292mod tests {
293 use iota_sdk_types::TransactionDigest;
294
295 use super::ObjectLocks;
296 use crate::execution_cache::{
297 ExecutionCacheWrite, writeback_cache::writeback_cache_tests::Scenario,
298 };
299
300 #[tokio::test]
301 async fn test_lock_conflict_does_not_count_against_pending_lock() {
302 telemetry_subscribers::init_for_testing();
303 Scenario::iterate(|mut s| async move {
304 s.with_created(&[1]);
305 s.do_tx().await;
306 let obj = s.obj_ref(1);
307
308 let holder = TransactionDigest::random();
309 let contender = TransactionDigest::random();
310
311 let locks = ObjectLocks::new();
314 locks
315 .try_set_transaction_lock(&obj, holder, &s.epoch_store)
316 .expect("first lock on a fresh object");
317 assert_eq!(*locks.locked_transactions.get(&obj).unwrap(), (1, holder));
318
319 locks
320 .try_set_transaction_lock(&obj, contender, &s.epoch_store)
321 .unwrap_err();
322 assert_eq!(
323 *locks.locked_transactions.get(&obj).unwrap(),
324 (1, holder),
325 "a rejected contender must not be counted as a holder"
326 );
327
328 locks
330 .try_set_transaction_lock(&obj, holder, &s.epoch_store)
331 .expect("same transaction may lock again");
332 assert_eq!(*locks.locked_transactions.get(&obj).unwrap(), (2, holder));
333 })
334 .await;
335 }
336
337 #[tokio::test]
338 async fn test_transaction_locks_are_exclusive() {
339 telemetry_subscribers::init_for_testing();
340 Scenario::iterate(|mut s| async move {
341 s.with_created(&[1, 2, 3]);
342 s.do_tx().await;
343
344 s.with_mutated(&[1, 2, 3]);
345 s.do_tx().await;
346
347 let new1 = s.obj_ref(1);
348 let new2 = s.obj_ref(2);
349 let new3 = s.obj_ref(3);
350
351 s.with_mutated(&[1, 2, 3]); let outputs = s.take_outputs();
353
354 let tx1 = s.make_signed_transaction(&outputs.transaction);
355
356 s.cache
357 .try_acquire_transaction_locks(&s.epoch_store, &[new1, new2], tx1)
358 .expect("locks should be available");
359 let cache = s.cache.clone();
362 let assert_no_staged_locks = || {
363 assert!(
364 cache
365 .object_locks_for_testing()
366 .locked_transactions
367 .is_empty()
368 )
369 };
370 assert_no_staged_locks();
371
372 s.with_created(&[4, 5]);
375 let tx2 = s.take_outputs().transaction.clone();
376 let tx2 = s.make_signed_transaction(&tx2);
377
378 s.cache
380 .try_acquire_transaction_locks(&s.epoch_store, &[new1, new2], tx2.clone())
381 .unwrap_err();
382 assert_no_staged_locks();
383
384 s.cache
386 .try_acquire_transaction_locks(&s.epoch_store, &[new3, new2], tx2.clone())
387 .unwrap_err();
388 assert_no_staged_locks();
389
390 s.cache
392 .try_acquire_transaction_locks(&s.epoch_store, &[new3], tx2)
393 .expect("new3 should be unlocked");
394 assert_no_staged_locks();
395 })
396 .await;
397 }
398
399 #[tokio::test]
400 async fn test_transaction_locks_are_durable() {
401 telemetry_subscribers::init_for_testing();
402 Scenario::iterate(|mut s| async move {
403 s.with_created(&[1, 2]);
404 s.do_tx().await;
405
406 let old2 = s.obj_ref(2);
407
408 s.with_mutated(&[1, 2]);
409 s.do_tx().await;
410
411 let new1 = s.obj_ref(1);
412 let new2 = s.obj_ref(2);
413
414 s.with_mutated(&[1, 2]); let outputs = s.take_outputs();
416
417 let tx = s.make_signed_transaction(&outputs.transaction);
418
419 s.cache
421 .try_acquire_transaction_locks(&s.epoch_store, &[new1, old2], tx.clone())
422 .unwrap_err();
423
424 s.cache
427 .try_acquire_transaction_locks(&s.epoch_store, &[new1, new2], tx)
428 .expect("new1 should be unlocked after revert");
429 })
430 .await;
431 }
432
433 #[tokio::test]
434 async fn test_acquire_transaction_locks_revert() {
435 telemetry_subscribers::init_for_testing();
436 Scenario::iterate(|mut s| async move {
437 s.with_created(&[1, 2]);
438 s.do_tx().await;
439
440 let old2 = s.obj_ref(2);
441
442 s.with_mutated(&[1, 2]);
443 s.do_tx().await;
444
445 let new1 = s.obj_ref(1);
446 let new2 = s.obj_ref(2);
447
448 s.with_mutated(&[1, 2]); let outputs = s.take_outputs();
450
451 let tx = s.make_signed_transaction(&outputs.transaction);
452
453 s.cache
455 .try_acquire_transaction_locks(&s.epoch_store, &[new1, old2], tx)
456 .unwrap_err();
457
458 s.with_created(&[4, 5]);
461 let tx2 = s.take_outputs().transaction.clone();
462 let tx2 = s.make_signed_transaction(&tx2);
463
464 s.cache
467 .try_acquire_transaction_locks(&s.epoch_store, &[new1, new2], tx2)
468 .expect("new1 should be unlocked after revert");
469 })
470 .await;
471 }
472
473 #[tokio::test]
474 async fn test_acquire_transaction_locks_is_sync() {
475 telemetry_subscribers::init_for_testing();
476 Scenario::iterate(|mut s| async move {
477 s.with_created(&[1, 2]);
478 s.do_tx().await;
479
480 let objects: Vec<_> = vec![s.object(1), s.object(2)]
481 .into_iter()
482 .map(|o| o.object_ref())
483 .collect();
484
485 s.with_mutated(&[1, 2]);
486 let outputs = s.take_outputs();
487
488 let tx2 = s.make_signed_transaction(&outputs.transaction);
489 s.cache
492 .try_acquire_transaction_locks(&s.epoch_store, &objects, tx2)
493 .unwrap();
494 })
495 .await;
496 }
497}