@ -61,10 +61,10 @@ Iterator* WritePreparedTxn::GetIterator(const ReadOptions& options,
Status WritePreparedTxn : : PrepareInternal ( ) {
Status WritePreparedTxn : : PrepareInternal ( ) {
WriteOptions write_options = write_options_ ;
WriteOptions write_options = write_options_ ;
write_options . disableWAL = false ;
write_options . disableWAL = false ;
const bool write_after_commit = true ;
const bool WRITE_AFTER_COMMIT = true ;
WriteBatchInternal : : MarkEndPrepare ( GetWriteBatch ( ) - > GetWriteBatch ( ) , name_ ,
WriteBatchInternal : : MarkEndPrepare ( GetWriteBatch ( ) - > GetWriteBatch ( ) , name_ ,
! write_after_commit ) ;
! WRITE_AFTER_COMMIT ) ;
const bool disable_memtable = true ;
const bool DISABLE_MEMTABLE = true ;
uint64_t seq_used = kMaxSequenceNumber ;
uint64_t seq_used = kMaxSequenceNumber ;
bool collapsed = GetWriteBatch ( ) - > Collapse ( ) ;
bool collapsed = GetWriteBatch ( ) - > Collapse ( ) ;
if ( collapsed ) {
if ( collapsed ) {
@ -74,7 +74,7 @@ Status WritePreparedTxn::PrepareInternal() {
Status s =
Status s =
db_impl_ - > WriteImpl ( write_options , GetWriteBatch ( ) - > GetWriteBatch ( ) ,
db_impl_ - > WriteImpl ( write_options , GetWriteBatch ( ) - > GetWriteBatch ( ) ,
/*callback*/ nullptr , & log_number_ , /*log ref*/ 0 ,
/*callback*/ nullptr , & log_number_ , /*log ref*/ 0 ,
! disable_memtable , & seq_used ) ;
! DISABLE_MEMTABLE , & seq_used ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
auto prepare_seq = seq_used ;
auto prepare_seq = seq_used ;
SetId ( prepare_seq ) ;
SetId ( prepare_seq ) ;
@ -91,34 +91,29 @@ Status WritePreparedTxn::CommitWithoutPrepareInternal() {
return CommitBatchInternal ( GetWriteBatch ( ) - > GetWriteBatch ( ) ) ;
return CommitBatchInternal ( GetWriteBatch ( ) - > GetWriteBatch ( ) ) ;
}
}
SequenceNumber WritePreparedTxn : : GetACommitSeqNumber ( SequenceNumber prep_seq ) {
if ( db_impl_ - > immutable_db_options ( ) . two_write_queues ) {
auto s = db_impl_ - > IncAndFetchSequenceNumber ( ) ;
# ifndef NDEBUG
MutexLock l ( & wpt_db_ - > seq_for_metadata_mutex_ ) ;
wpt_db_ - > seq_for_metadata . push_back ( s ) ;
# endif
return s ;
} else {
return prep_seq ;
}
}
Status WritePreparedTxn : : CommitBatchInternal ( WriteBatch * batch ) {
Status WritePreparedTxn : : CommitBatchInternal ( WriteBatch * batch ) {
// TODO(myabandeh): handle the duplicate keys in the batch
// TODO(myabandeh): handle the duplicate keys in the batch
// In the absence of Prepare markers, use Noop as a batch separator
// In the absence of Prepare markers, use Noop as a batch separator
WriteBatchInternal : : InsertNoop ( batch ) ;
WriteBatchInternal : : InsertNoop ( batch ) ;
const bool disable_memtable = true ;
const bool DISABLE_MEMTABLE = true ;
const uint64_t no_log_ref = 0 ;
const uint64_t no_log_ref = 0 ;
uint64_t seq_used = kMaxSequenceNumber ;
uint64_t seq_used = kMaxSequenceNumber ;
auto s = db_impl_ - > WriteImpl ( write_options_ , batch , nullptr , nullptr ,
auto s = db_impl_ - > WriteImpl ( write_options_ , batch , nullptr , nullptr ,
no_log_ref , ! disable_memtable , & seq_used ) ;
no_log_ref , ! DISABLE_MEMTABLE , & seq_used ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
uint64_t & prepare_seq = seq_used ;
uint64_t & prepare_seq = seq_used ;
uint64_t commit_seq = GetACommitSeqNumber ( prepare_seq ) ;
// Commit the batch by writing an empty batch to the queue that will release
// TODO(myabandeh): skip AddPrepared
// the commit sequence number to readers.
wpt_db_ - > AddPrepared ( prepare_seq ) ;
WritePreparedCommitEntryPreReleaseCallback update_commit_map (
wpt_db_ - > AddCommitted ( prepare_seq , commit_seq ) ;
wpt_db_ , db_impl_ , prepare_seq ) ;
WriteBatch empty_batch ;
empty_batch . PutLogData ( Slice ( ) ) ;
// In the absence of Prepare markers, use Noop as a batch separator
WriteBatchInternal : : InsertNoop ( & empty_batch ) ;
s = db_impl_ - > WriteImpl ( write_options_ , & empty_batch , nullptr , nullptr ,
no_log_ref , DISABLE_MEMTABLE , & seq_used ,
& update_commit_map ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
return s ;
return s ;
}
}
@ -137,28 +132,22 @@ Status WritePreparedTxn::CommitInternal() {
WriteBatchInternal : : SetAsLastestPersistentState ( working_batch ) ;
WriteBatchInternal : : SetAsLastestPersistentState ( working_batch ) ;
}
}
const bool disable_memtable = true ;
// TODO(myabandeh): Reject a commit request if AddCommitted cannot encode
// commit_seq. This happens if prep_seq <<< commit_seq.
auto prepare_seq = GetId ( ) ;
const bool includes_data = ! empty & & ! for_recovery ;
WritePreparedCommitEntryPreReleaseCallback update_commit_map (
wpt_db_ , db_impl_ , prepare_seq , includes_data ) ;
const bool disable_memtable = ! includes_data ;
uint64_t seq_used = kMaxSequenceNumber ;
uint64_t seq_used = kMaxSequenceNumber ;
// Since the prepared batch is directly written to memtable, there is already
// Since the prepared batch is directly written to memtable, there is already
// a connection between the memtable and its WAL, so there is no need to
// a connection between the memtable and its WAL, so there is no need to
// redundantly reference the log that contains the prepared data.
// redundantly reference the log that contains the prepared data.
const uint64_t zero_log_number = 0ull ;
const uint64_t zero_log_number = 0ull ;
auto s = db_impl_ - > WriteImpl (
auto s = db_impl_ - > WriteImpl ( write_options_ , working_batch , nullptr , nullptr ,
write_options_ , working_batch , nullptr , nullptr , zero_log_number ,
zero_log_number , disable_memtable , & seq_used ,
empty | | for_recovery ? disable_memtable : ! disable_memtable , & seq_used ) ;
& update_commit_map ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
uint64_t & commit_seq = seq_used ;
// TODO(myabandeh): Reject a commit request if AddCommitted cannot encode
// commit_seq. This happens if prep_seq <<< commit_seq.
auto prepare_seq = GetId ( ) ;
wpt_db_ - > AddCommitted ( prepare_seq , commit_seq ) ;
if ( ! empty & & ! for_recovery ) {
// Commit the data that is accompnaied with the commit marker
// TODO(myabandeh): skip AddPrepared
wpt_db_ - > AddPrepared ( commit_seq ) ;
uint64_t commit_seq_2 = GetACommitSeqNumber ( commit_seq ) ;
wpt_db_ - > AddCommitted ( commit_seq , commit_seq_2 ) ;
}
return s ;
return s ;
}
}
@ -232,19 +221,31 @@ Status WritePreparedTxn::RollbackInternal() {
}
}
// The Rollback marker will be used as a batch separator
// The Rollback marker will be used as a batch separator
WriteBatchInternal : : MarkRollback ( & rollback_batch , name_ ) ;
WriteBatchInternal : : MarkRollback ( & rollback_batch , name_ ) ;
const bool disable_memtable = true ;
const bool DISABLE_MEMTABLE = true ;
const uint64_t no_log_ref = 0 ;
const uint64_t no_log_ref = 0 ;
uint64_t seq_used = kMaxSequenceNumber ;
uint64_t seq_used = kMaxSequenceNumber ;
s = db_impl_ - > WriteImpl ( write_options_ , & rollback_batch , nullptr , nullptr ,
s = db_impl_ - > WriteImpl ( write_options_ , & rollback_batch , nullptr , nullptr ,
no_log_ref , ! disable_memtable , & seq_used ) ;
no_log_ref , ! DISABLE_MEMTABLE , & seq_used ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
if ( ! s . ok ( ) ) {
return s ;
}
uint64_t & prepare_seq = seq_used ;
uint64_t & prepare_seq = seq_used ;
uint64_t commit_seq = GetACommitSeqNumber ( prepare_seq ) ;
// Commit the batch by writing an empty batch to the queue that will release
// TODO(myabandeh): skip AddPrepared
// the commit sequence number to readers.
wpt_db_ - > AddPrepared ( prepare_seq ) ;
WritePreparedCommitEntryPreReleaseCallback update_commit_map (
wpt_db_ - > AddCommitted ( prepare_seq , commit_seq ) ;
wpt_db_ , db_impl_ , prepare_seq ) ;
WriteBatch empty_batch ;
empty_batch . PutLogData ( Slice ( ) ) ;
// In the absence of Prepare markers, use Noop as a batch separator
WriteBatchInternal : : InsertNoop ( & empty_batch ) ;
s = db_impl_ - > WriteImpl ( write_options_ , & empty_batch , nullptr , nullptr ,
no_log_ref , DISABLE_MEMTABLE , & seq_used ,
& update_commit_map ) ;
assert ( seq_used ! = kMaxSequenceNumber ) ;
// Mark the txn as rolled back
// Mark the txn as rolled back
wpt_db_ - > RollbackPrepared ( GetId ( ) , commit_seq ) ;
uint64_t & rollback_seq = seq_used ;
wpt_db_ - > RollbackPrepared ( GetId ( ) , rollback_seq ) ;
return s ;
return s ;
}
}