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