#pragma once #include "capnp/dynamic.h" #include #include #include #include "Python.h" class GILAcquire { public: GILAcquire() : gstate(PyGILState_Ensure()) {} ~GILAcquire() { PyGILState_Release(gstate); } PyGILState_STATE gstate; }; class GILRelease { public: GILRelease() { Py_UNBLOCK_THREADS } ~GILRelease() { Py_BLOCK_THREADS } PyThreadState *_save; // The macros above read/write from this variable }; class PyRefCounter { public: PyObject * obj; PyRefCounter(PyObject * o) : obj(o) { GILAcquire gil; Py_INCREF(obj); } PyRefCounter(const PyRefCounter & ref) : obj(ref.obj) { GILAcquire gil; Py_INCREF(obj); } ~PyRefCounter() { GILAcquire gil; Py_DECREF(obj); } }; inline kj::Own stealPyRef(PyObject* o) { auto ret = kj::heap(o); Py_DECREF(o); return ret; } ::kj::Promise> convert_to_pypromise(capnp::RemotePromise promise); inline ::kj::Promise> convert_to_pypromise(kj::Promise promise) { return promise.then([]() { GILAcquire gil; return kj::heap(Py_None); }); } void c_reraise_kj_exception(); void check_py_error(); ::kj::Promise> then(kj::Promise> promise, kj::Own func, kj::Own error_func); class PythonInterfaceDynamicImpl final: public capnp::DynamicCapability::Server { public: kj::Own py_server; kj::Own kj_loop; PythonInterfaceDynamicImpl(capnp::InterfaceSchema & schema, kj::Own _py_server, kj::Own kj_loop) : capnp::DynamicCapability::Server(schema), py_server(kj::mv(_py_server)), kj_loop(kj::mv(kj_loop)) { } ~PythonInterfaceDynamicImpl() { } kj::Promise call(capnp::InterfaceSchema::Method method, capnp::CallContext< capnp::DynamicStruct, capnp::DynamicStruct> context); }; class PyAsyncIoStream: public kj::AsyncIoStream { public: kj::Own protocol; PyAsyncIoStream(kj::Own protocol) : protocol(kj::mv(protocol)) {} ~PyAsyncIoStream(); kj::Promise tryRead(void* buffer, size_t minBytes, size_t maxBytes); kj::Promise write(const void* buffer, size_t size); kj::Promise write(kj::ArrayPtr> pieces); kj::Promise whenWriteDisconnected(); void shutdownWrite(); }; template inline void rejectDisconnected(kj::PromiseFulfiller& fulfiller, kj::StringPtr message) { fulfiller.reject(KJ_EXCEPTION(DISCONNECTED, message)); } inline void rejectVoidDisconnected(kj::PromiseFulfiller& fulfiller, kj::StringPtr message) { fulfiller.reject(KJ_EXCEPTION(DISCONNECTED, message)); } inline kj::Exception makeException(kj::StringPtr message) { return KJ_EXCEPTION(FAILED, message); } kj::Promise taskToPromise(kj::Own coroutine, PyObject* callback); ::kj::Promise> tryReadMessage(kj::AsyncIoStream& stream, capnp::ReaderOptions opts); void init_capnp_api();