118 CancelOptions options)
120 QMutexLocker locker(&m_mutex);
122 const auto oldState = state.loadRelaxed();
125 case QFutureInterfaceBase::CancelMode::CancelAndFinish:
126 if ((oldState & QFutureInterfaceBase::Finished)
127 && (oldState & QFutureInterfaceBase::Canceled)) {
130 switch_from_to(state, suspendingOrSuspended | QFutureInterfaceBase::Running,
131 QFutureInterfaceBase::Canceled | QFutureInterfaceBase::Finished);
133 case QFutureInterfaceBase::CancelMode::CancelOnly:
134 if (oldState & QFutureInterfaceBase::Canceled)
136 switch_from_to(state, suspendingOrSuspended, QFutureInterfaceBase::Canceled);
142 QMutexLocker continuationLocker(&continuationMutex);
145 QMutexLocker nextLocker(&next->continuationMutex);
146 if (next->continuationType == QFutureInterfaceBase::ContinuationType::Then) {
155 waitCondition.wakeAll();
156 pausedWaitCondition.wakeAll();
158 if (!(oldState & QFutureInterfaceBase::Canceled))
159 sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Canceled));
160 if (mode == QFutureInterfaceBase::CancelMode::CancelAndFinish
161 && !(oldState & QFutureInterfaceBase::Finished)) {
162 sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Finished));
193void QFutureInterfaceBase::cancelChain(QFutureInterfaceBase::CancelMode mode)
197 QMutexLocker locker(&d->continuationMutex);
198 QFutureInterfaceBasePrivate *prev = d->nonConcludedParent;
201 prev->cancelImpl(mode, QFutureInterfaceBasePrivate::CancelOption::None);
202 QMutexLocker prevLocker(&prev->continuationMutex);
203 prev = prev->nonConcludedParent;
207 d->cancelImpl(mode, QFutureInterfaceBasePrivate::CancelOption::CancelContinuations);
210void QFutureInterfaceBase::setSuspended(
bool suspend)
212 QMutexLocker locker(&d->m_mutex);
214 switch_on(d->state, Suspending);
215 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
217 switch_off(d->state, suspendingOrSuspended);
218 d->pausedWaitCondition.wakeAll();
219 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Resumed));
223void QFutureInterfaceBase::toggleSuspended()
225 QMutexLocker locker(&d->m_mutex);
226 if (d->state.loadRelaxed() & suspendingOrSuspended) {
227 switch_off(d->state, suspendingOrSuspended);
228 d->pausedWaitCondition.wakeAll();
229 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Resumed));
231 switch_on(d->state, Suspending);
232 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
236void QFutureInterfaceBase::reportSuspended()
const
241 QMutexLocker locker(&d->m_mutex);
242 const int state = d->state.loadRelaxed();
243 if (!(state & Suspending) || (state & Suspended))
246 switch_from_to(d->state, Suspending, Suspended);
247 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
250void QFutureInterfaceBase::setThrottled(
bool enable)
252 QMutexLocker lock(&d->m_mutex);
254 switch_on(d->state, Throttled);
256 switch_off(d->state, Throttled);
257 if (!(d->state.loadRelaxed() & suspendingOrSuspended))
258 d->pausedWaitCondition.wakeAll();
328void QFutureInterfaceBase::waitForResume()
332 const int state = d->state.loadRelaxed();
333 if (!(state & suspendingOrSuspended) || (state & Canceled))
337 QMutexLocker lock(&d->m_mutex);
338 const int state = d->state.loadRelaxed();
339 if (!(state & suspendingOrSuspended) || (state & Canceled))
343 const ThreadPoolThreadReleaser releaser(d->pool());
345 d->pausedWaitCondition.wait(&d->m_mutex);
348void QFutureInterfaceBase::suspendIfRequested()
350 const auto canSuspend = [] (
int state) {
352 return (state & suspendingOrSuspended) && !(state & Canceled);
357 const int state = d->state.loadRelaxed();
358 if (!canSuspend(state))
362 QMutexLocker lock(&d->m_mutex);
363 const int state = d->state.loadRelaxed();
364 if (!canSuspend(state))
368 if (!(state & Suspended)) {
370 switch_from_to(d->state, Suspending, Suspended);
371 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
375 const ThreadPoolThreadReleaser releaser(d->pool());
376 d->pausedWaitCondition.wait(&d->m_mutex);
415void QFutureInterfaceBase::reportStarted()
417 QMutexLocker locker(&d->m_mutex);
418 if (d->state.loadRelaxed() & (Started|Canceled|Finished))
420 d->setState(State(Started | Running));
421 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Started));
443void QFutureInterfaceBase::reportException(
const std::exception_ptr &exception)
446 QMutexLocker locker(&d->m_mutex);
447 if (d->state.loadRelaxed() & (Canceled|Finished))
450 d->hasException =
true;
451 d->data.setException(exception);
452 switch_on(d->state, Canceled);
453 d->waitCondition.wakeAll();
454 d->pausedWaitCondition.wakeAll();
455 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Canceled));
459void QFutureInterfaceBase::reportFinished()
461 QMutexLocker locker(&d->m_mutex);
463 switch_from_to(d->state, Running, Finished);
464 d->waitCondition.wakeAll();
465 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Finished));
494void QFutureInterfaceBase::waitForResult(
int resultIndex)
497 d->data.m_exceptionStore.rethrowException();
499 QMutexLocker lock(&d->m_mutex);
500 if (!isRunningOrPending())
506 d->pool()->d_func()->stealAndRunRunnable(d->runnable);
510 const int waitIndex = (resultIndex == -1) ? INT_MAX : resultIndex;
511 while (isRunningOrPending() && !d->internal_isResultReadyAt(waitIndex))
512 d->waitCondition.wait(&d->m_mutex);
515 d->data.m_exceptionStore.rethrowException();
518void QFutureInterfaceBase::waitForFinished()
520 QMutexLocker lock(&d->m_mutex);
521 const bool alreadyFinished = isFinished();
524 if (!alreadyFinished) {
525 d->pool()->d_func()->stealAndRunRunnable(d->runnable);
529 while (!isFinished())
530 d->waitCondition.wait(&d->m_mutex);
534 d->data.m_exceptionStore.rethrowException();
537void QFutureInterfaceBase::reportResultsReady(
int beginIndex,
int endIndex)
539 if (beginIndex == endIndex || (d->state.loadRelaxed() & (Canceled|Finished)))
542 d->waitCondition.wakeAll();
544 if (!d->m_progress) {
545 if (d->internal_updateProgressValue(d->m_progressValue + endIndex - beginIndex) ==
false) {
546 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady,
552 d->sendCallOuts(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
555 QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady,
560 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady, beginIndex, endIndex));
601void QFutureInterfaceBase::setProgressRange(
int minimum,
int maximum)
603 QMutexLocker locker(&d->m_mutex);
605 d->m_progress.reset(
new QFutureInterfaceBasePrivate::ProgressData());
606 d->m_progress->minimum = minimum;
607 d->m_progress->maximum = qMax(minimum, maximum);
608 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::ProgressRange, minimum, maximum));
609 d->m_progressValue = minimum;
623void QFutureInterfaceBase::setProgressValueAndText(
int progressValue,
624 const QString &progressText)
626 QMutexLocker locker(&d->m_mutex);
628 d->m_progress.reset(
new QFutureInterfaceBasePrivate::ProgressData());
630 const bool useProgressRange = (d->m_progress->maximum != 0) || (d->m_progress->minimum != 0);
632 && ((progressValue < d->m_progress->minimum) || (progressValue > d->m_progress->maximum))) {
636 if (d->m_progressValue >= progressValue)
639 if (d->state.loadRelaxed() & (Canceled|Finished))
642 if (d->internal_updateProgress(progressValue, progressText)) {
643 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
645 d->m_progress->text));
744 if (data.m_results.hasNextResult())
747 while ((state.loadRelaxed() & QFutureInterfaceBase::Running)
748 && data.m_results.hasNextResult() ==
false)
749 waitCondition.wait(&m_mutex);
751 return !(state.loadRelaxed() & QFutureInterfaceBase::Canceled)
752 && data.m_results.hasNextResult();
793 if ((enable && (state.loadRelaxed() & QFutureInterfaceBase::Throttled))
794 || (!enable && !(state.loadRelaxed() & QFutureInterfaceBase::Throttled)))
799 switch_on(state, QFutureInterfaceBase::Throttled);
801 switch_off(state, QFutureInterfaceBase::Throttled);
802 if (!(state.loadRelaxed() & suspendingOrSuspended))
803 pausedWaitCondition.wakeAll();
817 const QFutureCallOutEvent &callOutEvent2)
819 if (outputConnections.isEmpty())
822 for (
int i = 0; i < outputConnections.size(); ++i) {
823 QFutureCallOutInterface *iface = outputConnections.at(i);
824 iface->postCallOutEvent(callOutEvent1);
825 iface->postCallOutEvent(callOutEvent2);
834 QMutexLocker locker(&m_mutex);
836 const auto currentState = state.loadRelaxed();
837 if (currentState & QFutureInterfaceBase::Started) {
838 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Started));
840 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::ProgressRange,
842 m_progress->maximum));
843 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
847 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::ProgressRange,
850 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
857 QtPrivate::ResultIteratorBase it = data.m_results.begin();
858 while (it != data.m_results.end()) {
859 const int begin = it.resultIndex();
860 const int end = begin + it.batchSize();
861 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady,
868 if (currentState & QFutureInterfaceBase::Suspended)
869 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
870 else if (currentState & QFutureInterfaceBase::Suspending)
871 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
873 if (currentState & QFutureInterfaceBase::Canceled)
874 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Canceled));
876 if (currentState & QFutureInterfaceBase::Finished)
877 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Finished));
879 outputConnections.append(iface);
884 QMutexLocker lock(&m_mutex);
885 const qsizetype index = outputConnections.indexOf(iface);
888 outputConnections.removeAt(index);
890 iface->callOutInterfaceDisconnected();
898void QFutureInterfaceBase::setContinuation(std::function<
void (
const QFutureInterfaceBase &)> func,
899 void *continuationFutureData, ContinuationType type)
901 auto *futureData =
static_cast<QFutureInterfaceBasePrivate *>(continuationFutureData);
903 QMutexLocker lock(&d->continuationMutex);
908 if (d->continuationState != QFutureInterfaceBasePrivate::Cleaned) {
909 if (d->continuation) {
910 qWarning(
"Adding a continuation to a future which already has a continuation. "
911 "The existing continuation is overwritten.");
912 if (d->continuationData)
913 d->continuationData->nonConcludedParent =
nullptr;
918 Q_ASSERT_X(!futureData->nonConcludedParent,
"setContinuation",
919 "futureData already has a parent");
920 futureData->continuationType = type;
921 futureData->nonConcludedParent = d;
923 d->continuationData = futureData;
924 Q_ASSERT_X(!futureData || futureData->continuationType != ContinuationType::Unknown,
925 "setContinuation",
"Make sure to provide a correct continuation type!");
933 d->continuationExecuted =
true;
940 if (d->continuationState == QFutureInterfaceBasePrivate::Cleaned) {
944 futureData->nonConcludedParent =
nullptr;
946 d->continuation = std::move(func);
959void QFutureInterfaceBase::setContinuation(
const QObject *context, std::function<
void()> func,
960 const QVariant &continuationFuture,
961 ContinuationType type)
965 using FuncType =
void();
966 using Prototype =
typename QtPrivate::Callable<FuncType>::Function;
967 auto slotObj = QtPrivate::makeCallableObject<Prototype>(std::move(func));
969 auto slot = QtPrivate::SlotObjUniquePtr(slotObj);
971 auto *watcher =
new QObjectContinuationWrapper;
972 watcher->moveToThread(context->thread());
979 auto watcherMutex = std::make_shared<QRecursiveMutex>();
980 const auto destroyWatcher = [watcherMutex, watcher]()
mutable {
981 QMutexLocker lock(watcherMutex.get());
986 QObject::connect(watcher, &QObjectContinuationWrapper::run,
990 context, [slot = std::move(slot)] {
991 void *args[] = {
nullptr };
992 slot->call(
nullptr, args);
994 QObject::connect(watcher, &QObjectContinuationWrapper::run, watcher, destroyWatcher);
1001 QObject::connect(context, &QObject::destroyed, watcher, destroyWatcher);
1005 QFutureInterfaceBasePrivate *continuationFutureData =
nullptr;
1006 if (continuationFuture.isValid()) {
1007 Q_ASSERT(QLatin1StringView(continuationFuture.typeName())
1008 .startsWith(QLatin1StringView(
"QFutureInterface")));
1009 const auto continuationPtr =
1010 static_cast<
const QFutureInterfaceBase *>(continuationFuture.constData());
1011 continuationFutureData = continuationPtr->d;
1016 setContinuation([watcherMutex = std::move(watcherMutex),
1017 watcher = QPointer(watcher), continuationFuture]
1018 (
const QFutureInterfaceBase &parentData)
1020 Q_UNUSED(parentData);
1021 Q_UNUSED(continuationFuture);
1022 QMutexLocker lock(watcherMutex.get());
1024 emit watcher->run();
1025 }, continuationFutureData, type);
1039void QFutureInterfaceBase::runContinuation()
const
1046 bool ownsExtraRef =
false;
1047 QFutureInterfaceBasePrivate *dd = d;
1048 const auto derefGuard = qScopeGuard([&] {
1049 if (ownsExtraRef && !dd->refCount.deref())
1053 QMutexLocker lock(&dd->continuationMutex);
1054 if (dd->continuation && !dd->continuationExecuted) {
1057 if (dd->continuationData)
1058 dd->continuationData->nonConcludedParent =
nullptr;
1062 dd->continuationExecuted =
true;
1063 auto fn = std::move(dd->continuation);
1066 ownsExtraRef =
true;
1075 if (dd->continuationState != QFutureInterfaceBasePrivate::Cleaned)
1076 dd->continuation = std::move(fn);