Make ThreadSafeTransaction implementation less verbose
This commit is contained in:
parent
ec58f0f0cf
commit
c7d5bbd2e0
|
|
@ -133,10 +133,10 @@ ThreadSafeTransaction::ThreadSafeTransaction(DatabaseContext* cx) {
|
|||
// because the reference count of the DatabaseContext is solely managed from the main thread. If cx is destructed
|
||||
// immediately after this call, it will defer the DatabaseContext::delref (and onMainThread preserves the order of
|
||||
// these operations).
|
||||
ISingleThreadTransaction* tr = this->tr = ReadYourWritesTransaction::allocateOnForeignThread();
|
||||
this->tr = ReadYourWritesTransaction::allocateOnForeignThread();
|
||||
// No deferred error -- if the construction of the RYW transaction fails, we have no where to put it
|
||||
onMainThreadVoid(
|
||||
[tr, cx]() {
|
||||
[tr = this->tr, cx]() {
|
||||
cx->addref();
|
||||
new (tr) ReadYourWritesTransaction(Database(cx));
|
||||
},
|
||||
|
|
@ -150,18 +150,15 @@ ThreadSafeTransaction::~ThreadSafeTransaction() {
|
|||
}
|
||||
|
||||
void ThreadSafeTransaction::cancel() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr]() { tr->cancel(); }, nullptr);
|
||||
onMainThreadVoid([tr = this->tr]() { tr->cancel(); }, nullptr);
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::setVersion(Version v) {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, v]() { tr->setVersion(v); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, v]() { tr->setVersion(v); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
ThreadFuture<Version> ThreadSafeTransaction::getReadVersion() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr]() -> Future<Version> {
|
||||
return onMainThread([tr = this->tr]() -> Future<Version> {
|
||||
tr->checkDeferredError();
|
||||
return tr->getReadVersion();
|
||||
});
|
||||
|
|
@ -170,8 +167,7 @@ ThreadFuture<Version> ThreadSafeTransaction::getReadVersion() {
|
|||
ThreadFuture<Optional<Value>> ThreadSafeTransaction::get(const KeyRef& key, bool snapshot) {
|
||||
Key k = key;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, k, snapshot]() -> Future<Optional<Value>> {
|
||||
return onMainThread([tr = this->tr, k, snapshot]() -> Future<Optional<Value>> {
|
||||
tr->checkDeferredError();
|
||||
return tr->get(k, snapshot);
|
||||
});
|
||||
|
|
@ -180,8 +176,7 @@ ThreadFuture<Optional<Value>> ThreadSafeTransaction::get(const KeyRef& key, bool
|
|||
ThreadFuture<Key> ThreadSafeTransaction::getKey(const KeySelectorRef& key, bool snapshot) {
|
||||
KeySelector k = key;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, k, snapshot]() -> Future<Key> {
|
||||
return onMainThread([tr = this->tr, k, snapshot]() -> Future<Key> {
|
||||
tr->checkDeferredError();
|
||||
return tr->getKey(k, snapshot);
|
||||
});
|
||||
|
|
@ -190,8 +185,7 @@ ThreadFuture<Key> ThreadSafeTransaction::getKey(const KeySelectorRef& key, bool
|
|||
ThreadFuture<int64_t> ThreadSafeTransaction::getEstimatedRangeSizeBytes(const KeyRangeRef& keys) {
|
||||
KeyRange r = keys;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, r]() -> Future<int64_t> {
|
||||
return onMainThread([tr = this->tr, r]() -> Future<int64_t> {
|
||||
tr->checkDeferredError();
|
||||
return tr->getEstimatedRangeSizeBytes(r);
|
||||
});
|
||||
|
|
@ -201,8 +195,7 @@ ThreadFuture<Standalone<VectorRef<KeyRef>>> ThreadSafeTransaction::getRangeSplit
|
|||
int64_t chunkSize) {
|
||||
KeyRange r = range;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, r, chunkSize]() -> Future<Standalone<VectorRef<KeyRef>>> {
|
||||
return onMainThread([tr = this->tr, r, chunkSize]() -> Future<Standalone<VectorRef<KeyRef>>> {
|
||||
tr->checkDeferredError();
|
||||
return tr->getRangeSplitPoints(r, chunkSize);
|
||||
});
|
||||
|
|
@ -241,8 +234,7 @@ ThreadFuture<Standalone<RangeResultRef>> ThreadSafeTransaction::getRange(const K
|
|||
ThreadFuture<Standalone<VectorRef<const char*>>> ThreadSafeTransaction::getAddressesForKey(const KeyRef& key) {
|
||||
Key k = key;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, k]() -> Future<Standalone<VectorRef<const char*>>> {
|
||||
return onMainThread([tr = this->tr, k]() -> Future<Standalone<VectorRef<const char*>>> {
|
||||
tr->checkDeferredError();
|
||||
return tr->getAddressesForKey(k);
|
||||
});
|
||||
|
|
@ -251,21 +243,18 @@ ThreadFuture<Standalone<VectorRef<const char*>>> ThreadSafeTransaction::getAddre
|
|||
void ThreadSafeTransaction::addReadConflictRange(const KeyRangeRef& keys) {
|
||||
KeyRange r = keys;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, r]() { tr->addReadConflictRange(r); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, r]() { tr->addReadConflictRange(r); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::makeSelfConflicting() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr]() { tr->makeSelfConflicting(); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr]() { tr->makeSelfConflicting(); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::atomicOp(const KeyRef& key, const ValueRef& value, uint32_t operationType) {
|
||||
Key k = key;
|
||||
Value v = value;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, k, v, operationType]() { tr->atomicOp(k, v, operationType); },
|
||||
onMainThreadVoid([tr = this->tr, k, v, operationType]() { tr->atomicOp(k, v, operationType); },
|
||||
&tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
|
|
@ -273,24 +262,21 @@ void ThreadSafeTransaction::set(const KeyRef& key, const ValueRef& value) {
|
|||
Key k = key;
|
||||
Value v = value;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, k, v]() { tr->set(k, v); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, k, v]() { tr->set(k, v); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::clear(const KeyRangeRef& range) {
|
||||
KeyRange r = range;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, r]() { tr->clear(r); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, r]() { tr->clear(r); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::clear(const KeyRef& begin, const KeyRef& end) {
|
||||
Key b = begin;
|
||||
Key e = end;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid(
|
||||
[tr, b, e]() {
|
||||
[tr = this->tr, b, e]() {
|
||||
if (b > e)
|
||||
throw inverted_range();
|
||||
|
||||
|
|
@ -302,15 +288,13 @@ void ThreadSafeTransaction::clear(const KeyRef& begin, const KeyRef& end) {
|
|||
void ThreadSafeTransaction::clear(const KeyRef& key) {
|
||||
Key k = key;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, k]() { tr->clear(k); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, k]() { tr->clear(k); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
ThreadFuture<Void> ThreadSafeTransaction::watch(const KeyRef& key) {
|
||||
Key k = key;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, k]() -> Future<Void> {
|
||||
return onMainThread([tr = this->tr, k]() -> Future<Void> {
|
||||
tr->checkDeferredError();
|
||||
return tr->watch(k);
|
||||
});
|
||||
|
|
@ -319,13 +303,11 @@ ThreadFuture<Void> ThreadSafeTransaction::watch(const KeyRef& key) {
|
|||
void ThreadSafeTransaction::addWriteConflictRange(const KeyRangeRef& keys) {
|
||||
KeyRange r = keys;
|
||||
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr, r]() { tr->addWriteConflictRange(r); }, &tr->getMutableDeferredError());
|
||||
onMainThreadVoid([tr = this->tr, r]() { tr->addWriteConflictRange(r); }, &tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
ThreadFuture<Void> ThreadSafeTransaction::commit() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr]() -> Future<Void> {
|
||||
return onMainThread([tr = this->tr]() -> Future<Void> {
|
||||
tr->checkDeferredError();
|
||||
return tr->commit();
|
||||
});
|
||||
|
|
@ -337,13 +319,11 @@ Version ThreadSafeTransaction::getCommittedVersion() {
|
|||
}
|
||||
|
||||
ThreadFuture<int64_t> ThreadSafeTransaction::getApproximateSize() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr]() -> Future<int64_t> { return tr->getApproximateSize(); });
|
||||
return onMainThread([tr = this->tr]() -> Future<int64_t> { return tr->getApproximateSize(); });
|
||||
}
|
||||
|
||||
ThreadFuture<Standalone<StringRef>> ThreadSafeTransaction::getVersionstamp() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr]() -> Future<Standalone<StringRef>> { return tr->getVersionstamp(); });
|
||||
return onMainThread([tr = this->tr]() -> Future<Standalone<StringRef>> { return tr->getVersionstamp(); });
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::setOption(FDBTransactionOptions::Option option, Optional<StringRef> value) {
|
||||
|
|
@ -352,17 +332,15 @@ void ThreadSafeTransaction::setOption(FDBTransactionOptions::Option option, Opti
|
|||
TraceEvent("UnknownTransactionOption").detail("Option", option);
|
||||
throw invalid_option();
|
||||
}
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
Standalone<Optional<StringRef>> passValue = value;
|
||||
|
||||
// ThreadSafeTransaction is not allowed to do anything with options except pass them through to RYW.
|
||||
onMainThreadVoid([tr, option, passValue]() { tr->setOption(option, passValue.contents()); },
|
||||
onMainThreadVoid([tr = this->tr, option, passValue]() { tr->setOption(option, passValue.contents()); },
|
||||
&tr->getMutableDeferredError());
|
||||
}
|
||||
|
||||
ThreadFuture<Void> ThreadSafeTransaction::checkDeferredError() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr]() {
|
||||
return onMainThread([tr = this->tr]() {
|
||||
try {
|
||||
tr->checkDeferredError();
|
||||
} catch (Error& e) {
|
||||
|
|
@ -374,8 +352,7 @@ ThreadFuture<Void> ThreadSafeTransaction::checkDeferredError() {
|
|||
}
|
||||
|
||||
ThreadFuture<Void> ThreadSafeTransaction::onError(Error const& e) {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
return onMainThread([tr, e]() { return tr->onError(e); });
|
||||
return onMainThread([tr = this->tr, e]() { return tr->onError(e); });
|
||||
}
|
||||
|
||||
void ThreadSafeTransaction::operator=(ThreadSafeTransaction&& r) noexcept {
|
||||
|
|
@ -389,8 +366,7 @@ ThreadSafeTransaction::ThreadSafeTransaction(ThreadSafeTransaction&& r) noexcept
|
|||
}
|
||||
|
||||
void ThreadSafeTransaction::reset() {
|
||||
ISingleThreadTransaction* tr = this->tr;
|
||||
onMainThreadVoid([tr]() { tr->reset(); }, nullptr);
|
||||
onMainThreadVoid([tr = this->tr]() { tr->reset(); }, nullptr);
|
||||
}
|
||||
|
||||
extern const char* getSourceVersion();
|
||||
|
|
|
|||
Loading…
Reference in New Issue