mmdeploy/csrc/execution/run_loop.h
lzhangzz 46bfe0ac87
[Feature] New pipeline & executor for SDK (#497)
* executor prototype

* add split/when_all

* fix GCC build

* WIP let_value

* fix let_value

* WIP ensure_started

* ensure_started & start_detached

* fix let_value + when_all combo on MSVC 142

* fix static thread pool

* generic just, then, let_value, sync_wait

* minor

* generic split and when_all

* fully generic sender adapters

* when_all: workaround for GCC7

* support legacy spdlog

* fix memleak

* bulk

* static detector

* fix bulk & first pipeline

* bulk for static thread pools

* fix on MSVC

* WIP async batch submission

* WIP collation

* async batch

* fix detector

* fix async detector

* fix

* fix

* debug

* fix cuda allocator

* WIP type erased executor

* better type erasure

* simplify C API impl

* Expand & type erase TC

* deduction guide for type erased senders

* fix GCC build

* when_all for arrays of Value senders

* WIP pipeline v2

* WIP pipeline parser

* WIP timed batch operation

* add registry

* experiment

* fix pipeline

* naming

* fix mem-leak

* fix deferred batch operation

* WIP

* WIP configurable scheduler

* WIP configurable scheduler

* add comment

* parse scheduler config

* force link schedulers

* WIP pipeable sender

* WIP CPO

* ADL isolation and dismantle headers

* type erase single thread context

* fix MSVC build

* CPO

* replace decay_t with remove_cvref_t

* structure adjustment

* structure adjustment

* apply CPOs & C API rework

* refine C API

* detector async C API

* adjust detector async C API

* # Conflicts:
#	csrc/apis/c/detector.cpp

* fix when_all for type erased senders

* support void return for Then

* async detector

* fix some CPOs

* minor

* WIP rework capture mechanism for type erased types

* minor fix

* fix MSVC build

* move expand.h to execution

* make `Expand` pipeable

* fix type erased

* un-templatize `_TypeErasedOperation`

* re-work C API

* remove async_detector C API

* fix pipeline

* add flatten & unflatten

* fix flatten & unflatten

* add aync OCR demo

* config executor for nodes & better executor API

* working async OCR example

* minor

* dynamic batch via scheduler

* dynamic batch on `Value`

* fix MSVC build

* type erase dynamic batch scheduler

* sender as Python Awaitable

* naming

* naming

* add docs

* minor

* merge tmp branch

* unify C APIs

* fix ocr

* unify APIs

* fix typo

* update async OCR demo

* add v3 API text recognizer

* fix v3 API

* fix lint

* add license info & reformat

* add demo async_ocr_v2

* revert files

* revert files

* resolve link issues

* fix scheduler linkage for shared libs

* fix license header

* add docs for `mmdeploy_executor_split`

* add missing `mmdeploy_executor_transfer_just` and `mmdeploy_executor_execute`

* make `TimedSingleThreadContext` header only

* fix lint

* simplify type-erased sender
2022-06-01 14:10:43 +08:00

152 lines
3.4 KiB
C++

// Copyright (c) OpenMMLab. All rights reserved.
// Modified from
// https://github.com/brycelelbach/wg21_p2300_std_execution/blob/main/include/execution.hpp
#ifndef MMDEPLOY_CSRC_EXPERIMENTAL_EXECUTION_RUN_LOOP_H_
#define MMDEPLOY_CSRC_EXPERIMENTAL_EXECUTION_RUN_LOOP_H_
#include <condition_variable>
#include <mutex>
#include "utility.h"
namespace mmdeploy {
namespace __loop {
class RunLoop;
namespace __impl {
struct _Task {
virtual void _Execute() noexcept = 0;
_Task* next_ = nullptr;
};
template <typename Receiver>
struct _Operation {
struct type;
};
template <typename Receiver>
using operation_t = typename _Operation<remove_cvref_t<Receiver>>::type;
template <typename Receiver>
struct _Operation<Receiver>::type final : _Task {
friend void tag_invoke(start_t, type& op_state) noexcept { op_state._Start(); }
void _Execute() noexcept override { SetValue(std::move(receiver_)); }
void _Start() noexcept;
Receiver receiver_;
RunLoop* const loop_;
public:
template <class _Receiver2>
explicit type(_Receiver2&& receiver, RunLoop* loop)
: receiver_((_Receiver2 &&) receiver), loop_(loop) {}
};
} // namespace __impl
class RunLoop {
template <typename>
friend struct __impl::_Operation;
public:
class _Scheduler {
class _ScheduleTask {
friend _Scheduler;
template <typename Receiver>
friend __impl::operation_t<Receiver> tag_invoke(connect_t, const _ScheduleTask& self,
Receiver&& receiver) {
return __impl::operation_t<Receiver>{(Receiver &&) receiver, self.loop_};
}
RunLoop* const loop_;
public:
explicit _ScheduleTask(RunLoop* loop) noexcept : loop_(loop) {}
using value_types = std::tuple<>;
};
friend RunLoop;
explicit _Scheduler(RunLoop* loop) noexcept : loop_(loop) {}
public:
bool operator==(const _Scheduler& other) const noexcept { return loop_ == other.loop_; }
private:
friend _ScheduleTask tag_invoke(schedule_t, const _Scheduler& self) {
return _ScheduleTask{self.loop_};
}
RunLoop* loop_;
};
_Scheduler GetScheduler() { return _Scheduler{this}; }
void _Run();
void _Finish();
private:
void _push_back(__impl::_Task* task);
__impl::_Task* _pop_front();
std::mutex mutex_;
std::condition_variable cv_;
__impl::_Task* head_ = nullptr;
__impl::_Task* tail_ = nullptr;
bool stop_ = false;
};
namespace __impl {
template <typename Receiver>
inline void _Operation<Receiver>::type::_Start() noexcept {
loop_->_push_back(this);
}
} // namespace __impl
inline void RunLoop::_Run() {
while (auto* task = _pop_front()) {
task->_Execute();
}
}
inline void RunLoop::_Finish() {
std::lock_guard lock{mutex_};
stop_ = true;
cv_.notify_all();
}
inline void RunLoop::_push_back(__impl::_Task* task) {
std::lock_guard lock{mutex_};
if (head_ == nullptr) {
head_ = task;
} else {
tail_->next_ = task;
}
tail_ = task;
task->next_ = nullptr;
cv_.notify_one();
}
inline __impl::_Task* RunLoop::_pop_front() {
std::unique_lock lock{mutex_};
while (head_ == nullptr) {
if (stop_) {
return nullptr;
}
cv_.wait(lock);
}
auto* task = head_;
head_ = task->next_;
if (head_ == nullptr) {
tail_ = nullptr;
}
return task;
}
} // namespace __loop
using RunLoop = __loop::RunLoop;
} // namespace mmdeploy
#endif // MMDEPLOY_CSRC_EXPERIMENTAL_EXECUTION_RUN_LOOP_H_