ClickHouse/src/Common/AsyncTaskExecutor.cpp

122 lines
3.1 KiB
C++
Raw Normal View History

2023-03-03 19:30:43 +00:00
#include <Common/AsyncTaskExecutor.h>
namespace DB
{
AsyncTaskExecutor::AsyncTaskExecutor(std::unique_ptr<AsyncTask> task_) : task(std::move(task_))
{
createFiber();
}
void AsyncTaskExecutor::resume()
{
if (routine_is_finished)
return;
if (!checkBeforeTaskResume())
return;
{
std::lock_guard guard(fiber_lock);
if (is_cancelled)
return;
resumeUnlocked();
if (exception)
processException(exception);
}
afterTaskResume();
}
void AsyncTaskExecutor::resumeUnlocked()
{
fiber.resume();
2023-03-03 19:30:43 +00:00
}
void AsyncTaskExecutor::cancel()
{
std::lock_guard guard(fiber_lock);
is_cancelled = true;
2023-03-28 16:00:00 +00:00
cancelBefore();
2023-03-03 19:30:43 +00:00
destroyFiber();
2023-03-28 16:00:00 +00:00
cancelAfter();
2023-03-03 19:30:43 +00:00
}
void AsyncTaskExecutor::restart()
{
std::lock_guard guard(fiber_lock);
if (fiber)
destroyFiber();
createFiber();
routine_is_finished = false;
}
struct AsyncTaskExecutor::Routine
{
AsyncTaskExecutor & executor;
struct AsyncCallback
{
AsyncTaskExecutor & executor;
2023-05-22 18:22:05 +00:00
SuspendCallback suspend_callback;
2023-03-03 19:30:43 +00:00
void operator()(int fd, Poco::Timespan timeout, AsyncEventTimeoutType type, const std::string & desc, uint32_t events)
{
executor.processAsyncEvent(fd, timeout, type, desc, events);
2023-05-22 18:22:05 +00:00
suspend_callback();
2023-03-03 19:30:43 +00:00
executor.clearAsyncEvent();
}
};
2023-05-22 18:22:05 +00:00
void operator()(SuspendCallback suspend_callback)
2023-03-03 19:30:43 +00:00
{
2023-05-22 18:22:05 +00:00
auto async_callback = AsyncCallback{executor, suspend_callback};
2023-03-03 19:30:43 +00:00
try
{
executor.task->run(async_callback, suspend_callback);
}
catch (const boost::context::detail::forced_unwind &)
{
/// This exception is thrown by fiber implementation in case if fiber is being deleted but hasn't exited
/// It should not be caught or it will segfault.
/// Other exceptions must be caught
throw;
}
catch (...)
{
executor.exception = std::current_exception();
}
executor.routine_is_finished = true;
}
};
void AsyncTaskExecutor::createFiber()
{
fiber = Fiber(fiber_stack, Routine{*this});
2023-03-03 19:30:43 +00:00
}
void AsyncTaskExecutor::destroyFiber()
{
Fiber to_destroy = std::move(fiber);
2023-03-03 19:30:43 +00:00
}
String getSocketTimeoutExceededMessageByTimeoutType(AsyncEventTimeoutType type, Poco::Timespan timeout, const String & socket_description)
{
switch (type)
{
case AsyncEventTimeoutType::CONNECT:
return fmt::format("Timeout exceeded while connecting to socket ({}, connection timeout {} ms)", socket_description, timeout.totalMilliseconds());
case AsyncEventTimeoutType::RECEIVE:
return fmt::format("Timeout exceeded while reading from socket ({}, receive timeout {} ms)", socket_description, timeout.totalMilliseconds());
case AsyncEventTimeoutType::SEND:
return fmt::format("Timeout exceeded while writing to socket ({}, send timeout {} ms)", socket_description, timeout.totalMilliseconds());
default:
return fmt::format("Timeout exceeded while working with socket ({}, {} ms)", socket_description, timeout.totalMilliseconds());
}
}
}