/src/rocksdb/utilities/transactions/write_unprepared_txn_db.h
Line | Count | Source |
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 | | #include "utilities/transactions/write_prepared_txn_db.h" |
9 | | #include "utilities/transactions/write_unprepared_txn.h" |
10 | | |
11 | | namespace ROCKSDB_NAMESPACE { |
12 | | |
13 | | class WriteUnpreparedTxn; |
14 | | |
15 | | class WriteUnpreparedTxnDB : public WritePreparedTxnDB { |
16 | | public: |
17 | | using WritePreparedTxnDB::WritePreparedTxnDB; |
18 | | |
19 | | Status Initialize(const std::vector<size_t>& compaction_enabled_cf_indices, |
20 | | const std::vector<ColumnFamilyHandle*>& handles) override; |
21 | | |
22 | | Transaction* BeginTransaction(const WriteOptions& write_options, |
23 | | const TransactionOptions& txn_options, |
24 | | Transaction* old_txn) override; |
25 | | |
26 | | // Struct to hold ownership of snapshot and read callback for cleanup. |
27 | | struct IteratorState; |
28 | | |
29 | | using WritePreparedTxnDB::NewIterator; |
30 | | Iterator* NewIterator(const ReadOptions& _read_options, |
31 | | ColumnFamilyHandle* column_family, |
32 | | WriteUnpreparedTxn* txn); |
33 | | |
34 | | private: |
35 | | Status RollbackRecoveredTransaction(const DBImpl::RecoveredTransaction* rtxn); |
36 | | }; |
37 | | |
38 | | class WriteUnpreparedCommitEntryPreReleaseCallback : public PreReleaseCallback { |
39 | | // TODO(lth): Reduce code duplication with |
40 | | // WritePreparedCommitEntryPreReleaseCallback |
41 | | public: |
42 | | // includes_data indicates that the commit also writes non-empty |
43 | | // CommitTimeWriteBatch to memtable, which needs to be committed separately. |
44 | | WriteUnpreparedCommitEntryPreReleaseCallback( |
45 | | WritePreparedTxnDB* db, DBImpl* db_impl, |
46 | | const std::map<SequenceNumber, size_t>& unprep_seqs, |
47 | | size_t data_batch_cnt = 0, bool publish_seq = true) |
48 | 0 | : db_(db), |
49 | 0 | db_impl_(db_impl), |
50 | 0 | unprep_seqs_(unprep_seqs), |
51 | 0 | data_batch_cnt_(data_batch_cnt), |
52 | 0 | includes_data_(data_batch_cnt_ > 0), |
53 | 0 | publish_seq_(publish_seq) { |
54 | 0 | assert(unprep_seqs.size() > 0); |
55 | 0 | } |
56 | | |
57 | | Status Callback(SequenceNumber commit_seq, |
58 | | bool is_mem_disabled __attribute__((__unused__)), uint64_t, |
59 | 0 | size_t /*index*/, size_t /*total*/) override { |
60 | 0 | const uint64_t last_commit_seq = LIKELY(data_batch_cnt_ <= 1) |
61 | 0 | ? commit_seq |
62 | 0 | : commit_seq + data_batch_cnt_ - 1; |
63 | | // Recall that unprep_seqs maps (un)prepared_seq => prepare_batch_cnt. |
64 | 0 | for (const auto& s : unprep_seqs_) { |
65 | 0 | for (size_t i = 0; i < s.second; i++) { |
66 | 0 | db_->AddCommitted(s.first + i, last_commit_seq); |
67 | 0 | } |
68 | 0 | } |
69 | |
|
70 | 0 | if (includes_data_) { |
71 | 0 | assert(data_batch_cnt_); |
72 | | // Commit the data that is accompanied with the commit request |
73 | 0 | for (size_t i = 0; i < data_batch_cnt_; i++) { |
74 | | // For commit seq of each batch use the commit seq of the last batch. |
75 | | // This would make debugging easier by having all the batches having |
76 | | // the same sequence number. |
77 | 0 | db_->AddCommitted(commit_seq + i, last_commit_seq); |
78 | 0 | } |
79 | 0 | } |
80 | 0 | if (db_impl_->immutable_db_options().two_write_queues && publish_seq_) { |
81 | 0 | assert(is_mem_disabled); // implies the 2nd queue |
82 | | // Publish the sequence number. We can do that here assuming the callback |
83 | | // is invoked only from one write queue, which would guarantee that the |
84 | | // publish sequence numbers will be in order, i.e., once a seq is |
85 | | // published all the seq prior to that are also publishable. |
86 | 0 | db_impl_->SetLastPublishedSequence(last_commit_seq); |
87 | 0 | } |
88 | | // else SequenceNumber that is updated as part of the write already does the |
89 | | // publishing |
90 | 0 | return Status::OK(); |
91 | 0 | } |
92 | | |
93 | | private: |
94 | | WritePreparedTxnDB* db_; |
95 | | DBImpl* db_impl_; |
96 | | const std::map<SequenceNumber, size_t>& unprep_seqs_; |
97 | | size_t data_batch_cnt_; |
98 | | // Either because it is commit without prepare or it has a |
99 | | // CommitTimeWriteBatch |
100 | | bool includes_data_; |
101 | | // Should the callback also publishes the commit seq number |
102 | | bool publish_seq_; |
103 | | }; |
104 | | |
105 | | } // namespace ROCKSDB_NAMESPACE |