Skip to content

Commit

Permalink
fix race condition in scheduler reset and add test
Browse files Browse the repository at this point in the history
  • Loading branch information
csegarragonz committed Nov 23, 2021
1 parent 818a808 commit 4c0af62
Show file tree
Hide file tree
Showing 2 changed files with 33 additions and 4 deletions.
9 changes: 5 additions & 4 deletions src/scheduler/Executor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -61,15 +61,16 @@ void Executor::finish()

// Shut down thread pools and wait
for (int i = 0; i < threadPoolThreads.size(); i++) {
if (threadPoolThreads.at(i) == nullptr) {
continue;
}

// Send a kill message
SPDLOG_TRACE("Executor {} killing thread pool {}", id, i);
threadTaskQueues[i].enqueue(
ExecutorTask(POOL_SHUTDOWN, nullptr, nullptr, false, false));

// If already killed, move to the next thread
if (threadPoolThreads.at(i) == nullptr) {
continue;
}

// Await the thread
if (threadPoolThreads.at(i)->joinable()) {
threadPoolThreads.at(i)->join();
Expand Down
28 changes: 28 additions & 0 deletions tests/test/scheduler/test_executor.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -330,6 +330,34 @@ TEST_CASE_METHOD(TestExecutorFixture,
REQUIRE(restoreCount == 0);
}

TEST_CASE_METHOD(TestExecutorFixture,
"Test executing function repeatedly and flushing",
"[executor]")
{
std::shared_ptr<BatchExecuteRequest> req =
faabric::util::batchExecFactory("dummy", "simple", 1);
uint32_t msgId = req->messages().at(0).id();
std::vector<std::string> actualHosts;

int numRepeats = 20;
for (int i = 0; i < numRepeats; i++) {
std::vector<std::string> actualHosts =
executeWithTestExecutor(req, false);
faabric::Message result =
sch.getFunctionResult(msgId, SHORT_TEST_TIMEOUT_MS);
std::string expected =
fmt::format("Simple function {} executed", msgId);
REQUIRE(result.outputdata() == expected);

// We sleep for the same timeout threads have, to force a race condition
// between the scheduler's flush and the thread's own cleanup timeout
SLEEP_MS(conf.boundTimeout);

// Flush
sch.flushLocally();
}
}

TEST_CASE_METHOD(TestExecutorFixture,
"Test executing chained functions",
"[executor]")
Expand Down

0 comments on commit 4c0af62

Please sign in to comment.