Qt
Internal/Contributor docs for the Qt SDK. Note: These are NOT official API docs; those are found at https://doc.qt.io/
Loading...
Searching...
No Matches
qfutureinterface.cpp
Go to the documentation of this file.
1// Copyright (C) 2020 The Qt Company Ltd.
2// SPDX-License-Identifier: LicenseRef-Qt-Commercial OR LGPL-3.0-only OR GPL-2.0-only OR GPL-3.0-only
3// Qt-Security score:significant reason:default
4
5// qfutureinterface.h included from qfuture.h
6#include "qfuture.h"
8
9#include <QtCore/qatomic.h>
10#include <QtCore/qcoreapplication.h>
11#include <QtCore/qloggingcategory.h>
12#include <QtCore/qscopeguard.h>
13#include <QtCore/qthread.h>
14#include <QtCore/qvarlengtharray.h>
15#include <private/qthreadpool_p.h>
16#include <private/qobject_p.h>
17
18#include <climits> // For INT_MAX
19
20// GCC 12 gets confused about QFutureInterfaceBase::state, for some non-obvious
21// reason
22// warning: ‘unsigned int __atomic_or_fetch_4(volatile void*, unsigned int, int)’ writing 4 bytes into a region of size 0 overflows the destination [-Wstringop-overflow=]
23QT_WARNING_DISABLE_GCC("-Wstringop-overflow")
24
25QT_BEGIN_NAMESPACE
26
27Q_STATIC_LOGGING_CATEGORY(lcQFutureContinuations, "qt.core.qfuture.continuations")
28
29enum {
30 MaxProgressEmitsPerSecond = 25
31};
32
33namespace {
34class ThreadPoolThreadReleaser {
35 QThreadPool *m_pool;
36public:
37 Q_NODISCARD_CTOR
38 explicit ThreadPoolThreadReleaser(QThreadPool *pool)
39 : m_pool(pool)
40 { if (pool) pool->releaseThread(); }
41 ~ThreadPoolThreadReleaser()
42 { if (m_pool) m_pool->reserveThread(); }
43};
44
45const auto suspendingOrSuspended =
46 QFutureInterfaceBase::Suspending | QFutureInterfaceBase::Suspended;
47
48} // unnamed namespace
49
50namespace QtPrivate {
51
52void qfutureWarnIfUnusedResults(qsizetype numResults)
53{
54 if (numResults > 1) {
55 qCWarning(lcQFutureContinuations,
56 "Parent future has %" PRIdQSIZETYPE " result(s), but only the first result "
57 "will be handled in the continuation.",
58 numResults);
59 }
60}
61
62} // namespace QtPrivate
63
65{
67public:
70 {
71 }
72
74 void run();
75};
76
77QFutureCallOutInterface::~QFutureCallOutInterface()
78 = default;
79
81
82QFutureInterfaceBase::QFutureInterfaceBase(State initialState)
83 : d(new QFutureInterfaceBasePrivate(initialState))
84{ }
85
86QFutureInterfaceBase::QFutureInterfaceBase(const QFutureInterfaceBase &other)
87 : d(other.d)
88{
89 d->refCount.ref();
90}
91
92QFutureInterfaceBase::~QFutureInterfaceBase()
93{
94 if (d && !d->refCount.deref())
95 delete d;
96}
97
98static inline int switch_on(QAtomicInt &a, int which)
99{
100 return a.fetchAndOrRelaxed(which) | which;
101}
102
103static inline int switch_off(QAtomicInt &a, int which)
104{
105 return a.fetchAndAndRelaxed(~which) & ~which;
106}
107
108static inline int switch_from_to(QAtomicInt &a, int from, int to)
109{
110 const auto adjusted = [&](int old) { return (old & ~from) | to; };
111 int value = a.loadRelaxed();
112 while (!a.testAndSetRelaxed(value, adjusted(value), value))
113 qYieldCpu();
114 return value;
115}
116
117void QFutureInterfaceBasePrivate::cancelImpl(QFutureInterfaceBase::CancelMode mode,
118 CancelOptions options)
119{
120 QMutexLocker locker(&m_mutex);
121
122 const auto oldState = state.loadRelaxed();
123
124 switch (mode) {
125 case QFutureInterfaceBase::CancelMode::CancelAndFinish:
126 if ((oldState & QFutureInterfaceBase::Finished)
127 && (oldState & QFutureInterfaceBase::Canceled)) {
128 return;
129 }
130 switch_from_to(state, suspendingOrSuspended | QFutureInterfaceBase::Running,
131 QFutureInterfaceBase::Canceled | QFutureInterfaceBase::Finished);
132 break;
133 case QFutureInterfaceBase::CancelMode::CancelOnly:
134 if (oldState & QFutureInterfaceBase::Canceled)
135 return;
136 switch_from_to(state, suspendingOrSuspended, QFutureInterfaceBase::Canceled);
137 break;
138 }
139
140 if (options & CancelOption::CancelContinuations) {
141 // Cancel the continuations chain
142 QMutexLocker continuationLocker(&continuationMutex);
144 while (next) {
145 QMutexLocker nextLocker(&next->continuationMutex);
146 if (next->continuationType == QFutureInterfaceBase::ContinuationType::Then) {
147 next->continuationState = QFutureInterfaceBasePrivate::Canceled;
148 next = next->continuationData;
149 } else {
150 break;
151 }
152 }
153 }
154
155 waitCondition.wakeAll();
156 pausedWaitCondition.wakeAll();
157
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));
163 }
164
165 isValid = false;
166}
167
168void QFutureInterfaceBase::cancel()
169{
170 cancel(CancelMode::CancelOnly);
171}
172
173void QFutureInterfaceBase::cancelChain()
174{
175 cancelChain(CancelMode::CancelOnly);
176}
177
178void QFutureInterfaceBase::setAddResultsIfCanceledEnabled(bool enable)
179{
180 d->addResultsIfCanceled = enable;
181}
182
183bool QFutureInterfaceBase::isAddResultsIfCanceledEnabled() const
184{
185 return d->addResultsIfCanceled;
186}
187
188void QFutureInterfaceBase::cancel(QFutureInterfaceBase::CancelMode mode)
189{
190 d->cancelImpl(mode, QFutureInterfaceBasePrivate::CancelOption::CancelContinuations);
191}
192
193void QFutureInterfaceBase::cancelChain(QFutureInterfaceBase::CancelMode mode)
194{
195 // go up through the list of continuations, cancelling each of them
196 {
197 QMutexLocker locker(&d->continuationMutex);
198 QFutureInterfaceBasePrivate *prev = d->nonConcludedParent;
199 while (prev) {
200 // Do not cancel continuations, because we're going bottom-to-top
201 prev->cancelImpl(mode, QFutureInterfaceBasePrivate::CancelOption::None);
202 QMutexLocker prevLocker(&prev->continuationMutex);
203 prev = prev->nonConcludedParent;
204 }
205 }
206 // finally, cancel self and all next continuations
207 d->cancelImpl(mode, QFutureInterfaceBasePrivate::CancelOption::CancelContinuations);
208}
209
210void QFutureInterfaceBase::setSuspended(bool suspend)
211{
212 QMutexLocker locker(&d->m_mutex);
213 if (suspend) {
214 switch_on(d->state, Suspending);
215 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
216 } else {
217 switch_off(d->state, suspendingOrSuspended);
218 d->pausedWaitCondition.wakeAll();
219 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Resumed));
220 }
221}
222
223void QFutureInterfaceBase::toggleSuspended()
224{
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));
230 } else {
231 switch_on(d->state, Suspending);
232 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
233 }
234}
235
236void QFutureInterfaceBase::reportSuspended() const
237{
238 // Needs to be called when pause is in effect,
239 // i.e. no more events will be reported.
240
241 QMutexLocker locker(&d->m_mutex);
242 const int state = d->state.loadRelaxed();
243 if (!(state & Suspending) || (state & Suspended))
244 return;
245
246 switch_from_to(d->state, Suspending, Suspended);
247 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
248}
249
250void QFutureInterfaceBase::setThrottled(bool enable)
251{
252 QMutexLocker lock(&d->m_mutex);
253 if (enable) {
254 switch_on(d->state, Throttled);
255 } else {
256 switch_off(d->state, Throttled);
257 if (!(d->state.loadRelaxed() & suspendingOrSuspended))
258 d->pausedWaitCondition.wakeAll();
259 }
260}
261
262
263bool QFutureInterfaceBase::isRunning() const
264{
265 return queryState(Running);
266}
267
268bool QFutureInterfaceBase::isStarted() const
269{
270 return queryState(Started);
271}
272
273bool QFutureInterfaceBase::isCanceled() const
274{
275 return queryState(Canceled);
276}
277
278bool QFutureInterfaceBase::isFinished() const
279{
280 return queryState(Finished);
281}
282
283bool QFutureInterfaceBase::isSuspending() const
284{
285 return queryState(Suspending);
286}
287
288#if QT_DEPRECATED_SINCE(6, 0)
289bool QFutureInterfaceBase::isPaused() const
290{
291 return queryState(static_cast<State>(suspendingOrSuspended));
292}
293#endif
294
295bool QFutureInterfaceBase::isSuspended() const
296{
297 return queryState(Suspended);
298}
299
300bool QFutureInterfaceBase::isThrottled() const
301{
302 return queryState(Throttled);
303}
304
305bool QFutureInterfaceBase::isResultReadyAt(int index) const
306{
307 QMutexLocker lock(&d->m_mutex);
308 return d->internal_isResultReadyAt(index);
309}
310
311bool QFutureInterfaceBase::isValid() const
312{
313 const QMutexLocker lock(&d->m_mutex);
314 return d->isValid;
315}
316
317bool QFutureInterfaceBase::isRunningOrPending() const
318{
319 return queryState(static_cast<State>(Running | Pending));
320}
321
322bool QFutureInterfaceBase::waitForNextResult()
323{
324 QMutexLocker lock(&d->m_mutex);
325 return d->internal_waitForNextResult();
326}
327
328void QFutureInterfaceBase::waitForResume()
329{
330 // return early if possible to avoid taking the mutex lock.
331 {
332 const int state = d->state.loadRelaxed();
333 if (!(state & suspendingOrSuspended) || (state & Canceled))
334 return;
335 }
336
337 QMutexLocker lock(&d->m_mutex);
338 const int state = d->state.loadRelaxed();
339 if (!(state & suspendingOrSuspended) || (state & Canceled))
340 return;
341
342 // decrease active thread count since this thread will wait.
343 const ThreadPoolThreadReleaser releaser(d->pool());
344
345 d->pausedWaitCondition.wait(&d->m_mutex);
346}
347
348void QFutureInterfaceBase::suspendIfRequested()
349{
350 const auto canSuspend = [] (int state) {
351 // can suspend only if 1) in any suspend-related state; 2) not canceled
352 return (state & suspendingOrSuspended) && !(state & Canceled);
353 };
354
355 // return early if possible to avoid taking the mutex lock.
356 {
357 const int state = d->state.loadRelaxed();
358 if (!canSuspend(state))
359 return;
360 }
361
362 QMutexLocker lock(&d->m_mutex);
363 const int state = d->state.loadRelaxed();
364 if (!canSuspend(state))
365 return;
366
367 // Note: expecting that Suspending and Suspended are mutually exclusive
368 if (!(state & Suspended)) {
369 // switch state in case this is the first invocation
370 switch_from_to(d->state, Suspending, Suspended);
371 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
372 }
373
374 // decrease active thread count since this thread will wait.
375 const ThreadPoolThreadReleaser releaser(d->pool());
376 d->pausedWaitCondition.wait(&d->m_mutex);
377}
378
379int QFutureInterfaceBase::progressValue() const
380{
381 const QMutexLocker lock(&d->m_mutex);
382 return d->m_progressValue;
383}
384
385int QFutureInterfaceBase::progressMinimum() const
386{
387 const QMutexLocker lock(&d->m_mutex);
388 return d->m_progress ? d->m_progress->minimum : 0;
389}
390
391int QFutureInterfaceBase::progressMaximum() const
392{
393 const QMutexLocker lock(&d->m_mutex);
394 return d->m_progress ? d->m_progress->maximum : 0;
395}
396
397int QFutureInterfaceBase::resultCount() const
398{
399 QMutexLocker lock(&d->m_mutex);
400 return d->internal_resultCount();
401}
402
403QString QFutureInterfaceBase::progressText() const
404{
405 QMutexLocker locker(&d->m_mutex);
406 return d->m_progress ? d->m_progress->text : QString();
407}
408
409bool QFutureInterfaceBase::isProgressUpdateNeeded() const
410{
411 QMutexLocker locker(&d->m_mutex);
412 return !d->progressTime.isValid() || (d->progressTime.elapsed() > (1000 / MaxProgressEmitsPerSecond));
413}
414
415void QFutureInterfaceBase::reportStarted()
416{
417 QMutexLocker locker(&d->m_mutex);
418 if (d->state.loadRelaxed() & (Started|Canceled|Finished))
419 return;
420 d->setState(State(Started | Running));
421 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Started));
422 d->isValid = true;
423}
424
425void QFutureInterfaceBase::reportCanceled()
426{
427 cancel();
428}
429
430#ifndef QT_NO_EXCEPTIONS
431void QFutureInterfaceBase::reportException(const QException &exception)
432{
433 try {
434 exception.raise();
435 } catch (...) {
436 reportException(std::current_exception());
437 }
438}
439
440#if QT_VERSION < QT_VERSION_CHECK(7, 0, 0)
441void QFutureInterfaceBase::reportException(std::exception_ptr exception)
442#else
443void QFutureInterfaceBase::reportException(const std::exception_ptr &exception)
444#endif
445{
446 QMutexLocker locker(&d->m_mutex);
447 if (d->state.loadRelaxed() & (Canceled|Finished))
448 return;
449
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));
456}
457#endif
458
459void QFutureInterfaceBase::reportFinished()
460{
461 QMutexLocker locker(&d->m_mutex);
462 if (!isFinished()) {
463 switch_from_to(d->state, Running, Finished);
464 d->waitCondition.wakeAll();
465 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Finished));
466 }
467}
468
469void QFutureInterfaceBase::setExpectedResultCount(int resultCount)
470{
471 if (d->m_progress)
472 setProgressRange(0, resultCount);
473 d->m_expectedResultCount = resultCount;
474}
475
476int QFutureInterfaceBase::expectedResultCount()
477{
478 return d->m_expectedResultCount;
479}
480
481bool QFutureInterfaceBase::queryState(State state) const
482{
483 return d->state.loadRelaxed() & state;
484}
485
486int QFutureInterfaceBase::loadState() const
487{
488 // Used from ~QPromise, so this check is needed
489 if (!d)
490 return QFutureInterfaceBase::State::NoState;
491 return d->state.loadRelaxed();
492}
493
494void QFutureInterfaceBase::waitForResult(int resultIndex)
495{
496 if (d->hasException)
497 d->data.m_exceptionStore.rethrowException();
498
499 QMutexLocker lock(&d->m_mutex);
500 if (!isRunningOrPending())
501 return;
502 lock.unlock();
503
504 // To avoid deadlocks and reduce the number of threads used, try to
505 // run the runnable in the current thread.
506 d->pool()->d_func()->stealAndRunRunnable(d->runnable);
507
508 lock.relock();
509
510 const int waitIndex = (resultIndex == -1) ? INT_MAX : resultIndex;
511 while (isRunningOrPending() && !d->internal_isResultReadyAt(waitIndex))
512 d->waitCondition.wait(&d->m_mutex);
513
514 if (d->hasException)
515 d->data.m_exceptionStore.rethrowException();
516}
517
518void QFutureInterfaceBase::waitForFinished()
519{
520 QMutexLocker lock(&d->m_mutex);
521 const bool alreadyFinished = isFinished();
522 lock.unlock();
523
524 if (!alreadyFinished) {
525 d->pool()->d_func()->stealAndRunRunnable(d->runnable);
526
527 lock.relock();
528
529 while (!isFinished())
530 d->waitCondition.wait(&d->m_mutex);
531 }
532
533 if (d->hasException)
534 d->data.m_exceptionStore.rethrowException();
535}
536
537void QFutureInterfaceBase::reportResultsReady(int beginIndex, int endIndex)
538{
539 if (beginIndex == endIndex || (d->state.loadRelaxed() & (Canceled|Finished)))
540 return;
541
542 d->waitCondition.wakeAll();
543
544 if (!d->m_progress) {
545 if (d->internal_updateProgressValue(d->m_progressValue + endIndex - beginIndex) == false) {
546 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady,
547 beginIndex,
548 endIndex));
549 return;
550 }
551
552 d->sendCallOuts(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
553 d->m_progressValue,
554 QString()),
555 QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady,
556 beginIndex,
557 endIndex));
558 return;
559 }
560 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::ResultsReady, beginIndex, endIndex));
561}
562
563void QFutureInterfaceBase::setRunnable(QRunnable *runnable)
564{
565 d->runnable = runnable;
566}
567
568void QFutureInterfaceBase::setThreadPool(QThreadPool *pool)
569{
570 d->m_pool = pool;
571}
572
573QThreadPool *QFutureInterfaceBase::threadPool() const
574{
575 return d->m_pool;
576}
577
578void QFutureInterfaceBase::setFilterMode(bool enable)
579{
580 QMutexLocker locker(&d->m_mutex);
581 if (!hasException())
582 resultStoreBase().setFilterMode(enable);
583}
584
585/*!
586 \internal
587 Sets the progress range's minimum and maximum values to \a minimum and
588 \a maximum respectively.
589
590 If \a maximum is smaller than \a minimum, \a minimum becomes the only
591 legal value.
592
593 The progress value is reset to be \a minimum.
594
595 The progress range usage can be disabled by using setProgressRange(0, 0).
596 In this case progress value is also reset to 0.
597
598 The behavior of this method is mostly inspired by
599 \l QProgressBar::setRange.
600*/
601void QFutureInterfaceBase::setProgressRange(int minimum, int maximum)
602{
603 QMutexLocker locker(&d->m_mutex);
604 if (!d->m_progress)
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;
610}
611
612void QFutureInterfaceBase::setProgressValue(int progressValue)
613{
614 setProgressValueAndText(progressValue, QString());
615}
616
617/*!
618 \internal
619 In case of the \a progressValue falling out of the progress range,
620 this method has no effect.
621 Such behavior is inspired by \l QProgressBar::setValue.
622*/
623void QFutureInterfaceBase::setProgressValueAndText(int progressValue,
624 const QString &progressText)
625{
626 QMutexLocker locker(&d->m_mutex);
627 if (!d->m_progress)
628 d->m_progress.reset(new QFutureInterfaceBasePrivate::ProgressData());
629
630 const bool useProgressRange = (d->m_progress->maximum != 0) || (d->m_progress->minimum != 0);
631 if (useProgressRange
632 && ((progressValue < d->m_progress->minimum) || (progressValue > d->m_progress->maximum))) {
633 return;
634 }
635
636 if (d->m_progressValue >= progressValue)
637 return;
638
639 if (d->state.loadRelaxed() & (Canceled|Finished))
640 return;
641
642 if (d->internal_updateProgress(progressValue, progressText)) {
643 d->sendCallOut(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
644 d->m_progressValue,
645 d->m_progress->text));
646 }
647}
648
649QMutex &QFutureInterfaceBase::mutex() const
650{
651 return d->m_mutex;
652}
653
654bool QFutureInterfaceBase::hasException() const
655{
656 return d->hasException;
657}
658
659QtPrivate::ExceptionStore &QFutureInterfaceBase::exceptionStore()
660{
661 Q_ASSERT(d->hasException);
662 return d->data.m_exceptionStore;
663}
664
665QtPrivate::ResultStoreBase &QFutureInterfaceBase::resultStoreBase()
666{
667 Q_ASSERT(!d->hasException);
668 return d->data.m_results;
669}
670
671const QtPrivate::ResultStoreBase &QFutureInterfaceBase::resultStoreBase() const
672{
673 Q_ASSERT(!d->hasException);
674 return d->data.m_results;
675}
676
677QFutureInterfaceBase &QFutureInterfaceBase::operator=(const QFutureInterfaceBase &other)
678{
679 QFutureInterfaceBase copy(other);
680 swap(copy);
681 return *this;
682}
683
684// ### Qt 7: inline
685void QFutureInterfaceBase::swap(QFutureInterfaceBase &other) noexcept
686{
687 qSwap(d, other.d);
688}
689
690bool QFutureInterfaceBase::refT() const noexcept
691{
692 return d->refCount.refT();
693}
694
695bool QFutureInterfaceBase::derefT() const noexcept
696{
697 // Called from ~QFutureInterface
698 return !d || d->refCount.derefT();
699}
700
701void QFutureInterfaceBase::reset()
702{
703 d->m_progressValue = 0;
704 d->m_progress.reset();
705 d->progressTime.invalidate();
706 d->isValid = false;
707}
708
709void QFutureInterfaceBase::rethrowPossibleException()
710{
711 if (hasException())
712 exceptionStore().rethrowException();
713}
714
715QFutureInterfaceBasePrivate::QFutureInterfaceBasePrivate(QFutureInterfaceBase::State initialState)
717{
718 progressTime.invalidate();
719}
720
722{
723 if (hasException)
724 data.m_exceptionStore.~ExceptionStore();
725 else
726 data.m_results.~ResultStoreBase();
727}
728
730{
731 return hasException ? 0 : data.m_results.count(); // ### subtract canceled results.
732}
733
735{
736 return hasException ? false : (data.m_results.contains(index));
737}
738
740{
741 if (hasException)
742 return false;
743
744 if (data.m_results.hasNextResult())
745 return true;
746
747 while ((state.loadRelaxed() & QFutureInterfaceBase::Running)
748 && data.m_results.hasNextResult() == false)
749 waitCondition.wait(&m_mutex);
750
751 return !(state.loadRelaxed() & QFutureInterfaceBase::Canceled)
752 && data.m_results.hasNextResult();
753}
754
756{
757 if (m_progressValue >= progress)
758 return false;
759
760 m_progressValue = progress;
761
762 if (progressTime.isValid() && m_progressValue != 0) // make sure the first and last steps are emitted.
763 if (progressTime.elapsed() < (1000 / MaxProgressEmitsPerSecond))
764 return false;
765
766 progressTime.start();
767 return true;
768
769}
770
772 const QString &progressText)
773{
774 if (m_progressValue >= progress)
775 return false;
776
777 Q_ASSERT(m_progress);
778
779 m_progressValue = progress;
780 m_progress->text = progressText;
781
782 if (progressTime.isValid() && m_progressValue != m_progress->maximum) // make sure the first and last steps are emitted.
783 if (progressTime.elapsed() < (1000 / MaxProgressEmitsPerSecond))
784 return false;
785
786 progressTime.start();
787 return true;
788}
789
791{
792 // bail out if we are not changing the state
793 if ((enable && (state.loadRelaxed() & QFutureInterfaceBase::Throttled))
794 || (!enable && !(state.loadRelaxed() & QFutureInterfaceBase::Throttled)))
795 return;
796
797 // change the state
798 if (enable) {
799 switch_on(state, QFutureInterfaceBase::Throttled);
800 } else {
801 switch_off(state, QFutureInterfaceBase::Throttled);
802 if (!(state.loadRelaxed() & suspendingOrSuspended))
803 pausedWaitCondition.wakeAll();
804 }
805}
806
807void QFutureInterfaceBasePrivate::sendCallOut(const QFutureCallOutEvent &callOutEvent)
808{
809 if (outputConnections.isEmpty())
810 return;
811
812 for (int i = 0; i < outputConnections.size(); ++i)
813 outputConnections.at(i)->postCallOutEvent(callOutEvent);
814}
815
816void QFutureInterfaceBasePrivate::sendCallOuts(const QFutureCallOutEvent &callOutEvent1,
817 const QFutureCallOutEvent &callOutEvent2)
818{
819 if (outputConnections.isEmpty())
820 return;
821
822 for (int i = 0; i < outputConnections.size(); ++i) {
823 QFutureCallOutInterface *iface = outputConnections.at(i);
824 iface->postCallOutEvent(callOutEvent1);
825 iface->postCallOutEvent(callOutEvent2);
826 }
827}
828
829// This function connects an output interface (for example a QFutureWatcher)
830// to this future. While holding the lock we check the state and ready results
831// and add the appropriate callouts to the queue.
832void QFutureInterfaceBasePrivate::connectOutputInterface(QFutureCallOutInterface *iface)
833{
834 QMutexLocker locker(&m_mutex);
835
836 const auto currentState = state.loadRelaxed();
837 if (currentState & QFutureInterfaceBase::Started) {
838 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Started));
839 if (m_progress) {
840 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::ProgressRange,
841 m_progress->minimum,
842 m_progress->maximum));
843 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
844 m_progressValue,
845 m_progress->text));
846 } else {
847 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::ProgressRange,
848 0,
849 0));
850 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Progress,
851 m_progressValue,
852 QString()));
853 }
854 }
855
856 if (!hasException) {
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,
862 begin,
863 end));
864 it.batchedAdvance();
865 }
866 }
867
868 if (currentState & QFutureInterfaceBase::Suspended)
869 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Suspended));
870 else if (currentState & QFutureInterfaceBase::Suspending)
871 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Suspending));
872
873 if (currentState & QFutureInterfaceBase::Canceled)
874 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Canceled));
875
876 if (currentState & QFutureInterfaceBase::Finished)
877 iface->postCallOutEvent(QFutureCallOutEvent(QFutureCallOutEvent::Finished));
878
879 outputConnections.append(iface);
880}
881
882void QFutureInterfaceBasePrivate::disconnectOutputInterface(QFutureCallOutInterface *iface)
883{
884 QMutexLocker lock(&m_mutex);
885 const qsizetype index = outputConnections.indexOf(iface);
886 if (index == -1)
887 return;
888 outputConnections.removeAt(index);
889
890 iface->callOutInterfaceDisconnected();
891}
892
893void QFutureInterfaceBasePrivate::setState(QFutureInterfaceBase::State newState)
894{
895 state.storeRelaxed(newState);
896}
897
898void QFutureInterfaceBase::setContinuation(std::function<void (const QFutureInterfaceBase &)> func,
899 void *continuationFutureData, ContinuationType type)
900{
901 auto *futureData = static_cast<QFutureInterfaceBasePrivate *>(continuationFutureData);
902
903 QMutexLocker lock(&d->continuationMutex);
904
905 // Unless the continuation has been cleaned earlier, we have to
906 // store the move-only continuation, to guarantee that the associated
907 // future's data stays alive.
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;
914 }
915 if (futureData) {
916 // If we ever hit this, we should properly unlink the future before
917 // relinking to ourselves.
918 Q_ASSERT_X(!futureData->nonConcludedParent, "setContinuation",
919 "futureData already has a parent");
920 futureData->continuationType = type;
921 futureData->nonConcludedParent = d;
922 }
923 d->continuationData = futureData;
924 Q_ASSERT_X(!futureData || futureData->continuationType != ContinuationType::Unknown,
925 "setContinuation", "Make sure to provide a correct continuation type!");
926 }
927
928 // If the state is ready, run continuation immediately,
929 // otherwise save it for later.
930 // Ensure that we run the continuation after setting up the future chain links so that
931 // the continuation itself can modify the links again.
932 if (isFinished()) {
933 d->continuationExecuted = true;
934
935 lock.unlock();
936 func(*this);
937 lock.relock();
938 }
939
940 if (d->continuationState == QFutureInterfaceBasePrivate::Cleaned) {
941 // if the continuation that we've possibly run above has cleaned this
942 // future, we have to unlink ourselves from it again.
943 if (futureData)
944 futureData->nonConcludedParent = nullptr;
945 } else {
946 d->continuation = std::move(func);
947 }
948}
949
950/*
951 For continuations with context we expect all the needed data to be captured
952 directly by the continuation data, because this simplifies the slot
953 invocation. That's why func has no parameters.
954
955 We pass continuation data as a QVariant, because we need to keep the
956 QFutureInterface<T> for the entire lifetime of the continuation, but we
957 cannot pass a template type T as a parameter.
958*/
959void QFutureInterfaceBase::setContinuation(const QObject *context, std::function<void()> func,
960 const QVariant &continuationFuture,
961 ContinuationType type)
962{
963 Q_ASSERT(context);
964
965 using FuncType = void();
966 using Prototype = typename QtPrivate::Callable<FuncType>::Function;
967 auto slotObj = QtPrivate::makeCallableObject<Prototype>(std::move(func));
968
969 auto slot = QtPrivate::SlotObjUniquePtr(slotObj);
970
971 auto *watcher = new QObjectContinuationWrapper;
972 watcher->moveToThread(context->thread());
973
974 // We need to protect acccess to the watcher. The context object (and in turn, the watcher)
975 // could be destroyed while the continuation that emits the signal is running. We have to
976 // prevent that.
977 // The mutex has to be recursive, because the continuation itself could delete the context
978 // object (and thus the watcher), which will try to lock the mutex from the same thread twice.
979 auto watcherMutex = std::make_shared<QRecursiveMutex>();
980 const auto destroyWatcher = [watcherMutex, watcher]() mutable {
981 QMutexLocker lock(watcherMutex.get());
982 delete watcher;
983 };
984
985 // ### we're missing a convenient way to `QObject::connect()` to a `QSlotObjectBase`...
986 QObject::connect(watcher, &QObjectContinuationWrapper::run,
987 // for the following, cf. QMetaObject::invokeMethodImpl():
988 // we know `slot` is a lambda returning `void`, so we can just
989 // `call()` with `obj` and `args[0]` set to `nullptr`:
990 context, [slot = std::move(slot)] {
991 void *args[] = { nullptr }; // for `void` return value
992 slot->call(nullptr, args);
993 });
994 QObject::connect(watcher, &QObjectContinuationWrapper::run, watcher, destroyWatcher);
995
996 // We need to connect to destroyWatcher here, instead of delete or deleteLater().
997 // If the continuation is called from a separate thread, emit watcher->run() can't detect that
998 // the watcher has been deleted in the separate thread, causing a race condition and potential
999 // heap-use-after-free issue inside QObject::doActivate. destroyWatcher forces the deletion of
1000 // the watcher to occur after emit watcher->run() completes and prevents the race condition.
1001 QObject::connect(context, &QObject::destroyed, watcher, destroyWatcher);
1002
1003 // Extract a QFutureInterfaceBasePrivate pointer from the QVariant. We rely
1004 // on the fact that QVariant contains QFutureInterface<T>.
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;
1012 }
1013
1014 // Capture continuationFuture so that it lives as long as the continuation,
1015 // and the continuation data remains valid.
1016 setContinuation([watcherMutex = std::move(watcherMutex),
1017 watcher = QPointer(watcher), continuationFuture]
1018 (const QFutureInterfaceBase &parentData)
1019 {
1020 Q_UNUSED(parentData);
1021 Q_UNUSED(continuationFuture);
1022 QMutexLocker lock(watcherMutex.get());
1023 if (watcher)
1024 emit watcher->run();
1025 }, continuationFutureData, type);
1026}
1027
1028void QFutureInterfaceBase::cleanContinuation()
1029{
1030 if (!d)
1031 return;
1032
1033 QMutexLocker lock(&d->continuationMutex);
1034 d->continuation = nullptr;
1035 d->continuationState = QFutureInterfaceBasePrivate::Cleaned;
1036 d->continuationData = nullptr;
1037}
1038
1039void QFutureInterfaceBase::runContinuation() const
1040{
1041 // fn(*this) below runs arbitrary user code, which may (directly, or indirectly by
1042 // destrying the QFuture/QPromise whose runContinuation() is executing) - destroy
1043 // `this`.
1044 // So cache `d` in a local before calling fn(), and never touch `this` (or the
1045 // member `d`) again afterwards.
1046 bool ownsExtraRef = false;
1047 QFutureInterfaceBasePrivate *dd = d;
1048 const auto derefGuard = qScopeGuard([&] {
1049 if (ownsExtraRef && !dd->refCount.deref())
1050 delete dd;
1051 });
1052
1053 QMutexLocker lock(&dd->continuationMutex);
1054 if (dd->continuation && !dd->continuationExecuted) {
1055 // If we run the next continuation, then this future is concluded, so
1056 // we wouldn't need to revisit it in the cancelChain()
1057 if (dd->continuationData)
1058 dd->continuationData->nonConcludedParent = nullptr;
1059 // Save the continuation in a local function, to avoid calling
1060 // a null std::function below, in case cleanContinuation() is
1061 // called from some other thread right after unlock() below.
1062 dd->continuationExecuted = true;
1063 auto fn = std::move(dd->continuation);
1064
1065 dd->refCount.ref();
1066 ownsExtraRef = true;
1067
1068 lock.unlock();
1069 fn(*this);
1070
1071 lock.relock();
1072 // Unless the continuation has been cleaned earlier, we have to
1073 // store the move-only continuation, to guarantee that the associated
1074 // future's data stays alive.
1075 if (dd->continuationState != QFutureInterfaceBasePrivate::Cleaned)
1076 dd->continuation = std::move(fn);
1077 }
1078}
1079
1080bool QFutureInterfaceBase::isChainCanceled() const
1081{
1082 return isCanceled() || d->continuationState == QFutureInterfaceBasePrivate::Canceled;
1083}
1084
1085void QFutureInterfaceBase::setLaunchAsync(bool value)
1086{
1087 d->launchAsync = value;
1088}
1089
1090bool QFutureInterfaceBase::launchAsync() const
1091{
1092 return d->launchAsync;
1093}
1094
1095namespace QtFuture {
1096
1098{
1099 QFutureInterface<void> promise;
1100 promise.reportStarted();
1101 promise.reportFinished();
1102
1103 return promise.future();
1104}
1105
1106} // namespace QtFuture
1107
1108QT_END_NAMESPACE
1109
1110#include "qfutureinterface.moc"
bool internal_updateProgress(int progress, const QString &progressText=QString())
void sendCallOuts(const QFutureCallOutEvent &callOut1, const QFutureCallOutEvent &callOut2)
void sendCallOut(const QFutureCallOutEvent &callOut)
void setState(QFutureInterfaceBase::State state)
bool internal_updateProgressValue(int progress)
bool internal_isResultReadyAt(int index) const
void internal_setThrottled(bool enable)
QFutureInterfaceBasePrivate * continuationData
QFutureInterfaceBasePrivate(QFutureInterfaceBase::State initialState)
QFuture< void > makeReadyVoidFuture()
void qfutureWarnIfUnusedResults(qsizetype numResults)
static int switch_from_to(QAtomicInt &a, int from, int to)
static int switch_off(QAtomicInt &a, int which)
static int switch_on(QAtomicInt &a, int which)