[incfs] Allow multiple timed jobs at the same time point

Old code had a tiny chance of ignoring a job if it happens to
be scheduled to the exactly the same time as one already in the
queue. Not that it will ever happen, but better to fix it.

+ make the worker thread code slightly easier to reason about

Bug: 183243150
Test: atest IncrementalServiceTest
Change-Id: Ia3126d30e19edfd17f7c8da368e9763ca5501e84
This commit is contained in:
Yurii Zubrytskyi
2021-03-19 17:00:12 -07:00
parent 5fb05e5201
commit 0583e7f5c9

View File

@@ -255,7 +255,7 @@ public:
static JNIEnv* getOrAttachJniEnv(JavaVM* jvm); static JNIEnv* getOrAttachJniEnv(JavaVM* jvm);
class RealTimedQueueWrapper : public TimedQueueWrapper { class RealTimedQueueWrapper final : public TimedQueueWrapper {
public: public:
RealTimedQueueWrapper(JavaVM* jvm) { RealTimedQueueWrapper(JavaVM* jvm) {
mThread = std::thread([this, jvm]() { mThread = std::thread([this, jvm]() {
@@ -268,11 +268,11 @@ public:
CHECK(!mThread.joinable()) << "call stop first"; CHECK(!mThread.joinable()) << "call stop first";
} }
void addJob(MountId id, Milliseconds after, Job what) final { void addJob(MountId id, Milliseconds timeout, Job what) final {
const auto now = Clock::now(); const auto now = Clock::now();
{ {
std::unique_lock lock(mMutex); std::unique_lock lock(mMutex);
mJobs.insert(TimedJob{id, now + after, std::move(what)}); mJobs.insert(TimedJob{id, now + timeout, std::move(what)});
} }
mCondition.notify_all(); mCondition.notify_all();
} }
@@ -293,29 +293,28 @@ public:
private: private:
void runTimers() { void runTimers() {
static constexpr TimePoint kInfinityTs{Clock::duration::max()}; static constexpr TimePoint kInfinityTs{Clock::duration::max()};
TimePoint nextJobTs = kInfinityTs;
std::unique_lock lock(mMutex); std::unique_lock lock(mMutex);
for (;;) { for (;;) {
mCondition.wait_until(lock, nextJobTs, [this, nextJobTs]() { const TimePoint nextJobTs = mJobs.empty() ? kInfinityTs : mJobs.begin()->when;
mCondition.wait_until(lock, nextJobTs, [this, oldNextJobTs = nextJobTs]() {
const auto now = Clock::now(); const auto now = Clock::now();
const auto firstJobTs = !mJobs.empty() ? mJobs.begin()->when : kInfinityTs; const auto newFirstJobTs = !mJobs.empty() ? mJobs.begin()->when : kInfinityTs;
return !mRunning || firstJobTs <= now || firstJobTs < nextJobTs; return newFirstJobTs <= now || newFirstJobTs < oldNextJobTs || !mRunning;
}); });
if (!mRunning) { if (!mRunning) {
return; return;
} }
const auto now = Clock::now(); const auto now = Clock::now();
auto it = mJobs.begin(); // Always re-acquire begin(). We can't use it after unlock as mTimedJobs can change.
// Always acquire begin(). We can't use it after unlock as mTimedJobs can change. for (auto it = mJobs.begin(); it != mJobs.end() && it->when <= now;
for (; it != mJobs.end() && it->when <= now; it = mJobs.begin()) { it = mJobs.begin()) {
auto jobNode = mJobs.extract(it); auto jobNode = mJobs.extract(it);
lock.unlock(); lock.unlock();
jobNode.value().what(); jobNode.value().what();
lock.lock(); lock.lock();
} }
nextJobTs = it != mJobs.end() ? it->when : kInfinityTs;
} }
} }
@@ -328,7 +327,7 @@ private:
} }
}; };
bool mRunning = true; bool mRunning = true;
std::set<TimedJob> mJobs; std::multiset<TimedJob> mJobs;
std::condition_variable mCondition; std::condition_variable mCondition;
std::mutex mMutex; std::mutex mMutex;
std::thread mThread; std::thread mThread;