mmdeploy/csrc/apis/python/executor.cpp
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

70 lines
2.2 KiB
C++

// Copyright (c) OpenMMLab. All rights reserved.
#include "apis/python/common.h"
#include "core/utils/formatter.h"
#include "execution/execution.h"
#include "execution/schedulers/inlined_scheduler.h"
#include "execution/schedulers/registry.h"
#include "execution/schedulers/single_thread_context.h"
#include "execution/schedulers/static_thread_pool.h"
namespace mmdeploy {
namespace _python {
struct PySender {
TypeErasedSender<Value> sender_;
explicit PySender(TypeErasedSender<Value> sender) : sender_(std::move(sender)) {}
struct gil_guarded_deleter {
void operator()(py::object* p) const {
py::gil_scoped_acquire _;
delete p;
}
};
using object_ptr = std::unique_ptr<py::object, gil_guarded_deleter>;
py::object __await__() {
auto future = py::module::import("concurrent.futures").attr("Future")();
{
py::gil_scoped_release _;
StartDetached(std::move(sender_) |
Then([future = object_ptr{new py::object(future)}](const Value& value) mutable {
py::gil_scoped_acquire _;
future->attr("set_result")(ConvertToPyObject(value));
delete future.release();
}));
}
return py::module::import("asyncio").attr("wrap_future")(future).attr("__await__")();
}
};
} // namespace _python
using _python::PySender;
static void register_python_executor(py::module& m) {
py::class_<PySender, std::unique_ptr<PySender>>(m, "PySender")
.def("__await__", &PySender::__await__);
// m.def("test_async", [](const py::object& x) {
// static StaticThreadPool pool;
// TypeErasedScheduler<Value> scheduler{pool.GetScheduler()};
// auto sender = TransferJust(scheduler, ConvertToValue(x)) | Then([](Value x) {
// // std::this_thread::sleep_for(std::chrono::milliseconds(10));
// return Value(x.get<int>() * x.get<int>());
// });
// return std::make_unique<PySender>(std::move(sender));
// });
}
class PythonExecutorRegisterer {
public:
PythonExecutorRegisterer() { gPythonBindings().emplace("executor", register_python_executor); }
};
static PythonExecutorRegisterer python_executor_registerer;
} // namespace mmdeploy