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 #ifndef ROCKSDB_LITE
8 
9 #include <mutex>
10 #include <queue>
11 #include <set>
12 #include <string>
13 #include <unordered_map>
14 #include <vector>
15 
16 #include "db/db_iter.h"
17 #include "db/read_callback.h"
18 #include "db/snapshot_checker.h"
19 #include "rocksdb/db.h"
20 #include "rocksdb/options.h"
21 #include "rocksdb/utilities/transaction_db.h"
22 #include "util/cast_util.h"
23 #include "utilities/transactions/pessimistic_transaction.h"
24 #include "utilities/transactions/transaction_lock_mgr.h"
25 #include "utilities/transactions/write_prepared_txn.h"
26 
27 namespace ROCKSDB_NAMESPACE {
28 
29 class PessimisticTransactionDB : public TransactionDB {
30  public:
31   explicit PessimisticTransactionDB(DB* db,
32                                     const TransactionDBOptions& txn_db_options);
33 
34   explicit PessimisticTransactionDB(StackableDB* db,
35                                     const TransactionDBOptions& txn_db_options);
36 
37   virtual ~PessimisticTransactionDB();
38 
GetSnapshot()39   virtual const Snapshot* GetSnapshot() override { return db_->GetSnapshot(); }
40 
41   virtual Status Initialize(
42       const std::vector<size_t>& compaction_enabled_cf_indices,
43       const std::vector<ColumnFamilyHandle*>& handles);
44 
45   Transaction* BeginTransaction(const WriteOptions& write_options,
46                                 const TransactionOptions& txn_options,
47                                 Transaction* old_txn) override = 0;
48 
49   using StackableDB::Put;
50   virtual Status Put(const WriteOptions& options,
51                      ColumnFamilyHandle* column_family, const Slice& key,
52                      const Slice& val) override;
53 
54   using StackableDB::Delete;
55   virtual Status Delete(const WriteOptions& wopts,
56                         ColumnFamilyHandle* column_family,
57                         const Slice& key) override;
58 
59   using StackableDB::SingleDelete;
60   virtual Status SingleDelete(const WriteOptions& wopts,
61                               ColumnFamilyHandle* column_family,
62                               const Slice& key) override;
63 
64   using StackableDB::Merge;
65   virtual Status Merge(const WriteOptions& options,
66                        ColumnFamilyHandle* column_family, const Slice& key,
67                        const Slice& value) override;
68 
69   using TransactionDB::Write;
70   virtual Status Write(const WriteOptions& opts, WriteBatch* updates) override;
WriteWithConcurrencyControl(const WriteOptions & opts,WriteBatch * updates)71   inline Status WriteWithConcurrencyControl(const WriteOptions& opts,
72                                             WriteBatch* updates) {
73     // Need to lock all keys in this batch to prevent write conflicts with
74     // concurrent transactions.
75     Transaction* txn = BeginInternalTransaction(opts);
76     txn->DisableIndexing();
77 
78     auto txn_impl =
79         static_cast_with_check<PessimisticTransaction, Transaction>(txn);
80 
81     // Since commitBatch sorts the keys before locking, concurrent Write()
82     // operations will not cause a deadlock.
83     // In order to avoid a deadlock with a concurrent Transaction, Transactions
84     // should use a lock timeout.
85     Status s = txn_impl->CommitBatch(updates);
86 
87     delete txn;
88 
89     return s;
90   }
91 
92   using StackableDB::CreateColumnFamily;
93   virtual Status CreateColumnFamily(const ColumnFamilyOptions& options,
94                                     const std::string& column_family_name,
95                                     ColumnFamilyHandle** handle) override;
96 
97   using StackableDB::DropColumnFamily;
98   virtual Status DropColumnFamily(ColumnFamilyHandle* column_family) override;
99 
100   Status TryLock(PessimisticTransaction* txn, uint32_t cfh_id,
101                  const std::string& key, bool exclusive);
102 
103   void UnLock(PessimisticTransaction* txn, const TransactionKeyMap* keys);
104   void UnLock(PessimisticTransaction* txn, uint32_t cfh_id,
105               const std::string& key);
106 
107   void AddColumnFamily(const ColumnFamilyHandle* handle);
108 
109   static TransactionDBOptions ValidateTxnDBOptions(
110       const TransactionDBOptions& txn_db_options);
111 
GetTxnDBOptions()112   const TransactionDBOptions& GetTxnDBOptions() const {
113     return txn_db_options_;
114   }
115 
116   void InsertExpirableTransaction(TransactionID tx_id,
117                                   PessimisticTransaction* tx);
118   void RemoveExpirableTransaction(TransactionID tx_id);
119 
120   // If transaction is no longer available, locks can be stolen
121   // If transaction is available, try stealing locks directly from transaction
122   // It is the caller's responsibility to ensure that the referred transaction
123   // is expirable (GetExpirationTime() > 0) and that it is expired.
124   bool TryStealingExpiredTransactionLocks(TransactionID tx_id);
125 
126   Transaction* GetTransactionByName(const TransactionName& name) override;
127 
128   void RegisterTransaction(Transaction* txn);
129   void UnregisterTransaction(Transaction* txn);
130 
131   // not thread safe. current use case is during recovery (single thread)
132   void GetAllPreparedTransactions(std::vector<Transaction*>* trans) override;
133 
134   TransactionLockMgr::LockStatusData GetLockStatusData() override;
135 
136   std::vector<DeadlockPath> GetDeadlockInfoBuffer() override;
137   void SetDeadlockInfoBufferSize(uint32_t target_size) override;
138 
139   // The default implementation does nothing. The actual implementation is moved
140   // to the child classes that actually need this information. This was due to
141   // an odd performance drop we observed when the added std::atomic member to
142   // the base class even when the subclass do not read it in the fast path.
UpdateCFComparatorMap(const std::vector<ColumnFamilyHandle * > &)143   virtual void UpdateCFComparatorMap(const std::vector<ColumnFamilyHandle*>&) {}
UpdateCFComparatorMap(ColumnFamilyHandle *)144   virtual void UpdateCFComparatorMap(ColumnFamilyHandle*) {}
145 
146  protected:
147   DBImpl* db_impl_;
148   std::shared_ptr<Logger> info_log_;
149   const TransactionDBOptions txn_db_options_;
150 
151   void ReinitializeTransaction(
152       Transaction* txn, const WriteOptions& write_options,
153       const TransactionOptions& txn_options = TransactionOptions());
154 
155   virtual Status VerifyCFOptions(const ColumnFamilyOptions& cf_options);
156 
157  private:
158   friend class WritePreparedTxnDB;
159   friend class WritePreparedTxnDBMock;
160   friend class WriteUnpreparedTxn;
161   friend class TransactionTest_DoubleCrashInRecovery_Test;
162   friend class TransactionTest_DoubleEmptyWrite_Test;
163   friend class TransactionTest_DuplicateKeys_Test;
164   friend class TransactionTest_PersistentTwoPhaseTransactionTest_Test;
165   friend class TransactionTest_TwoPhaseDoubleRecoveryTest_Test;
166   friend class TransactionTest_TwoPhaseOutOfOrderDelete_Test;
167   friend class TransactionStressTest_TwoPhaseLongPrepareTest_Test;
168   friend class WriteUnpreparedTransactionTest_RecoveryTest_Test;
169   friend class WriteUnpreparedTransactionTest_MarkLogWithPrepSection_Test;
170   TransactionLockMgr lock_mgr_;
171 
172   // Must be held when adding/dropping column families.
173   InstrumentedMutex column_family_mutex_;
174   Transaction* BeginInternalTransaction(const WriteOptions& options);
175 
176   // Used to ensure that no locks are stolen from an expirable transaction
177   // that has started a commit. Only transactions with an expiration time
178   // should be in this map.
179   std::mutex map_mutex_;
180   std::unordered_map<TransactionID, PessimisticTransaction*>
181       expirable_transactions_map_;
182 
183   // map from name to two phase transaction instance
184   std::mutex name_map_mutex_;
185   std::unordered_map<TransactionName, Transaction*> transactions_;
186 
187   // Signal that we are testing a crash scenario. Some asserts could be relaxed
188   // in such cases.
TEST_Crash()189   virtual void TEST_Crash() {}
190 };
191 
192 // A PessimisticTransactionDB that writes the data to the DB after the commit.
193 // In this way the DB only contains the committed data.
194 class WriteCommittedTxnDB : public PessimisticTransactionDB {
195  public:
WriteCommittedTxnDB(DB * db,const TransactionDBOptions & txn_db_options)196   explicit WriteCommittedTxnDB(DB* db,
197                                const TransactionDBOptions& txn_db_options)
198       : PessimisticTransactionDB(db, txn_db_options) {}
199 
WriteCommittedTxnDB(StackableDB * db,const TransactionDBOptions & txn_db_options)200   explicit WriteCommittedTxnDB(StackableDB* db,
201                                const TransactionDBOptions& txn_db_options)
202       : PessimisticTransactionDB(db, txn_db_options) {}
203 
~WriteCommittedTxnDB()204   virtual ~WriteCommittedTxnDB() {}
205 
206   Transaction* BeginTransaction(const WriteOptions& write_options,
207                                 const TransactionOptions& txn_options,
208                                 Transaction* old_txn) override;
209 
210   // Optimized version of ::Write that makes use of skip_concurrency_control
211   // hint
212   using TransactionDB::Write;
213   virtual Status Write(const WriteOptions& opts,
214                        const TransactionDBWriteOptimizations& optimizations,
215                        WriteBatch* updates) override;
216   virtual Status Write(const WriteOptions& opts, WriteBatch* updates) override;
217 };
218 
219 }  // namespace ROCKSDB_NAMESPACE
220 #endif  // ROCKSDB_LITE
221