| 147 | } |
| 148 | |
| 149 | StatelessTaskExecutor::TaskStatus StatelessTaskExecutor::getStatus(const String & task_id, UInt64 wait_milliseconds) |
| 150 | { |
| 151 | /// Make a copy of task completion future to wait for it outside of the lock |
| 152 | std::shared_future<String> completion_future; |
| 153 | std::shared_ptr<Progress> progress; |
| 154 | { |
| 155 | std::lock_guard lock(tasks_mutex); |
| 156 | auto it = tasks.find(task_id); |
| 157 | if (it == tasks.end()) |
| 158 | return TaskStatus{Result::UnknownTaskId, "", {}}; |
| 159 | completion_future = it->second->completion_future; |
| 160 | progress = it->second->progress; |
| 161 | } |
| 162 | |
| 163 | if (completion_future.valid() && completion_future.wait_for(std::chrono::milliseconds(wait_milliseconds)) == std::future_status::timeout) |
| 164 | { |
| 165 | Progress progress_delta = progress->fetchAndResetPiecewiseAtomically(); |
| 166 | return TaskStatus{Result::TaskRunnig, "", std::move(progress_delta)}; |
| 167 | } |
| 168 | |
| 169 | Progress progress_delta = progress->fetchAndResetPiecewiseAtomically(); |
| 170 | auto error_message = completion_future.get(); |
| 171 | if (error_message.empty()) |
| 172 | return TaskStatus{Result::TaskFinished, "", std::move(progress_delta)}; |
| 173 | else |
| 174 | return TaskStatus{Result::TaskFailed, error_message, std::move(progress_delta)}; |
| 175 | } |
| 176 | |
| 177 | StatelessTaskExecutor::Result StatelessTaskExecutor::cancelTask(const String & task_id) |
| 178 | { |
no test coverage detected