1 // Copyright (c) 2011-present, Facebook, Inc. All rights reserved. 2 // This source code is licensed under both the GPLv2 (found in the 3 // COPYING file in the root directory) and Apache 2.0 License 4 // (found in the LICENSE.Apache file in the root directory). 5 6 #pragma once 7 8 #ifndef ROCKSDB_LITE 9 10 #include <algorithm> 11 #include <atomic> 12 #include <mutex> 13 #include <stack> 14 #include <string> 15 #include <unordered_map> 16 #include <vector> 17 18 #include "db/write_callback.h" 19 #include "rocksdb/db.h" 20 #include "rocksdb/slice.h" 21 #include "rocksdb/snapshot.h" 22 #include "rocksdb/status.h" 23 #include "rocksdb/types.h" 24 #include "rocksdb/utilities/transaction.h" 25 #include "rocksdb/utilities/transaction_db.h" 26 #include "rocksdb/utilities/write_batch_with_index.h" 27 #include "util/autovector.h" 28 #include "utilities/transactions/transaction_base.h" 29 #include "utilities/transactions/transaction_util.h" 30 31 namespace ROCKSDB_NAMESPACE { 32 33 class PessimisticTransactionDB; 34 35 // A transaction under pessimistic concurrency control. This class implements 36 // the locking API and interfaces with the lock manager as well as the 37 // pessimistic transactional db. 38 class PessimisticTransaction : public TransactionBaseImpl { 39 public: 40 PessimisticTransaction(TransactionDB* db, const WriteOptions& write_options, 41 const TransactionOptions& txn_options, 42 const bool init = true); 43 // No copying allowed 44 PessimisticTransaction(const PessimisticTransaction&) = delete; 45 void operator=(const PessimisticTransaction&) = delete; 46 47 virtual ~PessimisticTransaction(); 48 49 void Reinitialize(TransactionDB* txn_db, const WriteOptions& write_options, 50 const TransactionOptions& txn_options); 51 52 Status Prepare() override; 53 54 Status Commit() override; 55 56 // It is basically Commit without going through Prepare phase. The write batch 57 // is also directly provided instead of expecting txn to gradually batch the 58 // transactions writes to an internal write batch. 59 Status CommitBatch(WriteBatch* batch); 60 61 Status Rollback() override; 62 63 Status RollbackToSavePoint() override; 64 65 Status SetName(const TransactionName& name) override; 66 67 // Generate a new unique transaction identifier 68 static TransactionID GenTxnID(); 69 GetID()70 TransactionID GetID() const override { return txn_id_; } 71 GetWaitingTxns(uint32_t * column_family_id,std::string * key)72 std::vector<TransactionID> GetWaitingTxns(uint32_t* column_family_id, 73 std::string* key) const override { 74 std::lock_guard<std::mutex> lock(wait_mutex_); 75 std::vector<TransactionID> ids(waiting_txn_ids_.size()); 76 if (key) *key = waiting_key_ ? *waiting_key_ : ""; 77 if (column_family_id) *column_family_id = waiting_cf_id_; 78 std::copy(waiting_txn_ids_.begin(), waiting_txn_ids_.end(), ids.begin()); 79 return ids; 80 } 81 SetWaitingTxn(autovector<TransactionID> ids,uint32_t column_family_id,const std::string * key)82 void SetWaitingTxn(autovector<TransactionID> ids, uint32_t column_family_id, 83 const std::string* key) { 84 std::lock_guard<std::mutex> lock(wait_mutex_); 85 waiting_txn_ids_ = ids; 86 waiting_cf_id_ = column_family_id; 87 waiting_key_ = key; 88 } 89 ClearWaitingTxn()90 void ClearWaitingTxn() { 91 std::lock_guard<std::mutex> lock(wait_mutex_); 92 waiting_txn_ids_.clear(); 93 waiting_cf_id_ = 0; 94 waiting_key_ = nullptr; 95 } 96 97 // Returns the time (in microseconds according to Env->GetMicros()) 98 // that this transaction will be expired. Returns 0 if this transaction does 99 // not expire. GetExpirationTime()100 uint64_t GetExpirationTime() const { return expiration_time_; } 101 102 // returns true if this transaction has an expiration_time and has expired. 103 bool IsExpired() const; 104 105 // Returns the number of microseconds a transaction can wait on acquiring a 106 // lock or -1 if there is no timeout. GetLockTimeout()107 int64_t GetLockTimeout() const { return lock_timeout_; } SetLockTimeout(int64_t timeout)108 void SetLockTimeout(int64_t timeout) override { 109 lock_timeout_ = timeout * 1000; 110 } 111 112 // Returns true if locks were stolen successfully, false otherwise. 113 bool TryStealingLocks(); 114 IsDeadlockDetect()115 bool IsDeadlockDetect() const override { return deadlock_detect_; } 116 GetDeadlockDetectDepth()117 int64_t GetDeadlockDetectDepth() const { return deadlock_detect_depth_; } 118 119 protected: 120 // Refer to 121 // TransactionOptions::use_only_the_last_commit_time_batch_for_recovery 122 bool use_only_the_last_commit_time_batch_for_recovery_ = false; 123 124 virtual Status PrepareInternal() = 0; 125 126 virtual Status CommitWithoutPrepareInternal() = 0; 127 128 // batch_cnt if non-zero is the number of sub-batches. A sub-batch is a batch 129 // with no duplicate keys. If zero, then the number of sub-batches is unknown. 130 virtual Status CommitBatchInternal(WriteBatch* batch, 131 size_t batch_cnt = 0) = 0; 132 133 virtual Status CommitInternal() = 0; 134 135 virtual Status RollbackInternal() = 0; 136 137 virtual void Initialize(const TransactionOptions& txn_options); 138 139 Status LockBatch(WriteBatch* batch, TransactionKeyMap* keys_to_unlock); 140 141 Status TryLock(ColumnFamilyHandle* column_family, const Slice& key, 142 bool read_only, bool exclusive, const bool do_validate = true, 143 const bool assume_tracked = false) override; 144 145 void Clear() override; 146 147 PessimisticTransactionDB* txn_db_impl_; 148 DBImpl* db_impl_; 149 150 // If non-zero, this transaction should not be committed after this time (in 151 // microseconds according to Env->NowMicros()) 152 uint64_t expiration_time_; 153 154 private: 155 friend class TransactionTest_ValidateSnapshotTest_Test; 156 // Used to create unique ids for transactions. 157 static std::atomic<TransactionID> txn_id_counter_; 158 159 // Unique ID for this transaction 160 TransactionID txn_id_; 161 162 // IDs for the transactions that are blocking the current transaction. 163 // 164 // empty if current transaction is not waiting. 165 autovector<TransactionID> waiting_txn_ids_; 166 167 // The following two represents the (cf, key) that a transaction is waiting 168 // on. 169 // 170 // If waiting_key_ is not null, then the pointer should always point to 171 // a valid string object. The reason is that it is only non-null when the 172 // transaction is blocked in the TransactionLockMgr::AcquireWithTimeout 173 // function. At that point, the key string object is one of the function 174 // parameters. 175 uint32_t waiting_cf_id_; 176 const std::string* waiting_key_; 177 178 // Mutex protecting waiting_txn_ids_, waiting_cf_id_ and waiting_key_. 179 mutable std::mutex wait_mutex_; 180 181 // Timeout in microseconds when locking a key or -1 if there is no timeout. 182 int64_t lock_timeout_; 183 184 // Whether to perform deadlock detection or not. 185 bool deadlock_detect_; 186 187 // Whether to perform deadlock detection or not. 188 int64_t deadlock_detect_depth_; 189 190 // Refer to TransactionOptions::skip_concurrency_control 191 bool skip_concurrency_control_; 192 193 virtual Status ValidateSnapshot(ColumnFamilyHandle* column_family, 194 const Slice& key, 195 SequenceNumber* tracked_at_seq); 196 197 void UnlockGetForUpdate(ColumnFamilyHandle* column_family, 198 const Slice& key) override; 199 }; 200 201 class WriteCommittedTxn : public PessimisticTransaction { 202 public: 203 WriteCommittedTxn(TransactionDB* db, const WriteOptions& write_options, 204 const TransactionOptions& txn_options); 205 // No copying allowed 206 WriteCommittedTxn(const WriteCommittedTxn&) = delete; 207 void operator=(const WriteCommittedTxn&) = delete; 208 ~WriteCommittedTxn()209 virtual ~WriteCommittedTxn() {} 210 211 private: 212 Status PrepareInternal() override; 213 214 Status CommitWithoutPrepareInternal() override; 215 216 Status CommitBatchInternal(WriteBatch* batch, size_t batch_cnt) override; 217 218 Status CommitInternal() override; 219 220 Status RollbackInternal() override; 221 }; 222 223 } // namespace ROCKSDB_NAMESPACE 224 225 #endif // ROCKSDB_LITE 226