Skip to main content

iota_core/execution_cache/
object_locks.rs

1// Copyright (c) Mysten Labs, Inc.
2// Modifications Copyright (c) 2024 IOTA Stiftung
3// SPDX-License-Identifier: Apache-2.0
4
5use 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    // When acquire transaction locks, lock entries are briefly inserted into this map. The map
23    // exists to provide atomic test-and-set operations on the locks. After all locks have been
24    // inserted into the map, they are written to the db, and then all locks are removed from
25    // the map.
26    //
27    // After a transaction has been executed, newly created objects are available to be locked.
28    // But, because of crash recovery, we cannot rule out that a lock may already exist in the db
29    // for those objects. Therefore we do a db read for each object we are locking.
30    //
31    // TODO: find a strategy to allow us to avoid db reads for each object.
32    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        // We don't consult the in-memory state here. We are only interested in state
48        // that has been committed to the db. This is because in memory state is
49        // reverted if the transaction is not successfully locked.
50        epoch_store.tables()?.get_locked_transaction(obj_ref)
51    }
52
53    /// Attempts to atomically test-and-set a transaction lock on an object.
54    /// If the lock is already set to a conflicting transaction, an error is
55    /// returned and the cached state is left untouched. If the lock is not
56    /// set, or is already set to the same transaction, the lock is set and the
57    /// caller must release it with `clear_cached_locks`.
58    fn try_set_transaction_lock(
59        &self,
60        obj_ref: &ObjectReference,
61        new_lock: LockDetails,
62        epoch_store: &AuthorityPerEpochStore,
63    ) -> IotaResult {
64        // entry holds a lock on the dashmap shard, so this function operates atomically
65        let entry = self.locked_transactions.entry(*obj_ref);
66
67        // TODO: currently, the common case for this code is that we will miss the cache
68        // and read from the db. It is difficult to implement negative caching, since we
69        // may have restarted, in which case there could be locks in the db that we do
70        // not have in the cache. We may want to explore strategies for proving there
71        // cannot be a lock in the db that we do not know about. Two possibilities are:
72        //
73        // 1. Read all locks into memory at startup (and keep them there). The lifetime
74        //    of locks is relatively short in the common case, so this might be
75        //    feasible.
76        // 2. Find some strategy to distinguish between the cases where we are
77        //    re-executing old transactions after restarting vs executing transactions
78        //    that we have never seen before. The output objects of novel transactions
79        //    cannot previously have been locked on this validator.
80        //
81        // Solving this is not terribly important as it is not in the execution path,
82        // and hence only improves the latency of transaction signing, not
83        // transaction execution
84        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        // A rejected caller never calls `clear_cached_locks` for this object, so
96        // the entry must only be created or counted once the lock is known to be
97        // compatible.
98        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                // this is impossible because the only case in which we overwrite a
173                // lock is when the lock is from a previous epoch. but we are holding
174                // execution_lock, so the epoch cannot have changed.
175                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    /// Validates that all owned input objects exist and their versions/digests
202    /// match the live objects. Does not acquire any locks.
203    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        // Only live objects can be locked
235        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        // Sort the objects before locking. This is not required by the protocol (since
243        // it's okay to reject any equivocating tx). However, this does prevent
244        // a confusing error on the client. Consider the case:
245        //   TX1: [o1, o2];
246        //   TX2: [o2, o1];
247        // If two threads race to acquire these locks, they might both acquire the first
248        // object, then error when trying to acquire the second. The error
249        // returned to the client would say that there is a conflicting tx on
250        // that object, but in fact neither object was locked and the tx was never
251        // signed. If one client then retries, they will succeed (counterintuitively).
252        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        // Note that this function does not have to operate atomically. If there are two
259        // racing threads, then they are either trying to lock the same
260        // transaction (in which case both will succeed), or they are trying to
261        // lock the same object in two different transactions, in which case the
262        // sender has equivocated, and we are under no obligation to help them form a
263        // cert.
264        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                    // revert all pending writes and return error
269                    // Note that reverting is not required for liveness, since a well formed and
270                    // un-equivocating txn cannot fail to acquire locks.
271                    // However, reverting is easy enough to do in this implementation that we do it
272                    // anyway.
273                    self.clear_cached_locks(&locks_to_write);
274                    return Err(e);
275                }
276            }
277        }
278
279        // commit all writes to DB
280        epoch_store
281            .tables()?
282            .write_transaction_locks(transaction, locks_to_write.iter().cloned())?;
283
284        // remove pending locks from unbounded storage
285        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            // A signer of `holder` is mid-flight: its entry is set but not yet
312            // written to the db and released.
313            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            // Another signer of the same transaction still shares the entry.
329            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]); // begin forming a tx but never execute it
352            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            // Entries are staged only until the locks are written, and a
360            // rejected conflict must leave the map exactly as it found it.
361            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            // this tx doesn't use the actual objects in question, but we just need
373            // something to insert into the table.
374            s.with_created(&[4, 5]);
375            let tx2 = s.take_outputs().transaction.clone();
376            let tx2 = s.make_signed_transaction(&tx2);
377
378            // both locks are held by tx1, so this should fail
379            s.cache
380                .try_acquire_transaction_locks(&s.epoch_store, &[new1, new2], tx2.clone())
381                .unwrap_err();
382            assert_no_staged_locks();
383
384            // new3 is lockable, but new2 is not, so this should fail
385            s.cache
386                .try_acquire_transaction_locks(&s.epoch_store, &[new3, new2], tx2.clone())
387                .unwrap_err();
388            assert_no_staged_locks();
389
390            // new3 is unlocked
391            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]); // begin forming a tx but never execute it
415            let outputs = s.take_outputs();
416
417            let tx = s.make_signed_transaction(&outputs.transaction);
418
419            // fails because we are referring to an old object
420            s.cache
421                .try_acquire_transaction_locks(&s.epoch_store, &[new1, old2], tx.clone())
422                .unwrap_err();
423
424            // succeeds because the above call releases the lock on new1 after failing
425            // to get the lock on old2
426            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]); // begin forming a tx but never execute it
449            let outputs = s.take_outputs();
450
451            let tx = s.make_signed_transaction(&outputs.transaction);
452
453            // fails because we are referring to an old object
454            s.cache
455                .try_acquire_transaction_locks(&s.epoch_store, &[new1, old2], tx)
456                .unwrap_err();
457
458            // this tx doesn't use the actual objects in question, but we just need
459            // something to insert into the table.
460            s.with_created(&[4, 5]);
461            let tx2 = s.take_outputs().transaction.clone();
462            let tx2 = s.make_signed_transaction(&tx2);
463
464            // succeeds because the above call releases the lock on new1 after failing
465            // to get the lock on old2
466            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            // assert that acquire_transaction_locks is sync in non-simtest, which causes
490            // the fail_point_async! macros above to be elided
491            s.cache
492                .try_acquire_transaction_locks(&s.epoch_store, &objects, tx2)
493                .unwrap();
494        })
495        .await;
496    }
497}