Get rid of VoidPromise and (almost) Promise
We only retain RemotePromise for its pipelining capabilities.
This commit is contained in:
@@ -10,7 +10,7 @@
|
||||
cimport cython # noqa: E402
|
||||
|
||||
from capnp.helpers.helpers cimport init_capnp_api
|
||||
from capnp.includes.capnp_cpp cimport AsyncIoStream, WaitScope, PyPromise, VoidPromise, EventPort, EventLoop, WaitScope, LowLevelAsyncIoProvider, AsyncIoProvider, newAsyncIoProvider, Canceler, PyAsyncIoStream, PromiseFulfiller, VoidPromiseFulfiller, tryReadMessage, writeMessage, makeException
|
||||
from capnp.includes.capnp_cpp cimport AsyncIoStream, WaitScope, PyPromise, VoidPromise, EventPort, EventLoop, Canceler, PyAsyncIoStream, PromiseFulfiller, VoidPromiseFulfiller, tryReadMessage, writeMessage, makeException
|
||||
from capnp.includes.schema_cpp cimport (MessageReader,)
|
||||
|
||||
from cpython cimport array, Py_buffer, PyObject_CheckBuffer, memoryview, buffer
|
||||
@@ -60,12 +60,7 @@ def deregister_all_types():
|
||||
|
||||
# By making it public, we'll be able to call it from capabilityHelper.h
|
||||
cdef api object wrap_dynamic_struct_reader(Response & r) with gil:
|
||||
return _Response()._init_childptr(new Response(moveResponse(r)), None)
|
||||
|
||||
|
||||
cdef api object wrap_remote_call(object func, Response & r):
|
||||
response = _Response()._init_childptr(new Response(moveResponse(r)), None)
|
||||
return func(response)
|
||||
return _Response()._init_childptr(new Response(move(r)), None)
|
||||
|
||||
cdef _find_field_order(struct_node):
|
||||
return [f.name for f in sorted(struct_node.fields, key=_attrgetter('codeOrder'))]
|
||||
@@ -185,7 +180,7 @@ cdef class _KjExceptionWrapper:
|
||||
cdef capnp.Exception * thisptr
|
||||
|
||||
cdef _init(self, capnp.Exception & other):
|
||||
self.thisptr = new capnp.Exception(moveException(other))
|
||||
self.thisptr = new capnp.Exception(move(other))
|
||||
return self
|
||||
|
||||
def __dealloc__(self):
|
||||
@@ -306,12 +301,6 @@ ctypedef fused _DynamicSetterClasses:
|
||||
DynamicStruct_Builder
|
||||
|
||||
|
||||
ctypedef fused PromiseTypes:
|
||||
_Promise
|
||||
_RemotePromise
|
||||
_VoidPromise
|
||||
|
||||
|
||||
cdef extern from "Python.h":
|
||||
cdef int PyObject_GetBuffer(object, Py_buffer *view, int flags)
|
||||
cdef void PyBuffer_Release(Py_buffer *view)
|
||||
@@ -327,20 +316,6 @@ cdef extern from "capnp/list.h" namespace " ::capnp":
|
||||
uint size()
|
||||
|
||||
|
||||
cdef extern from "<utility>" namespace "std":
|
||||
C_DynamicStruct.Pipeline moveStructPipeline"std::move"(C_DynamicStruct.Pipeline)
|
||||
C_DynamicOrphan moveOrphan"std::move"(C_DynamicOrphan)
|
||||
Request moveRequest"std::move"(Request)
|
||||
Response moveResponse"std::move"(Response)
|
||||
PyPromise movePromise"std::move"(PyPromise)
|
||||
VoidPromise moveVoidPromise"std::move"(VoidPromise)
|
||||
RemotePromise moveRemotePromise"std::move"(RemotePromise)
|
||||
CallContext moveCallContext"std::move"(CallContext)
|
||||
Own[AsyncIoStream] moveOwnAsyncIOStream"std::move"(Own[AsyncIoStream])
|
||||
capnp.Exception moveException"std::move"(capnp.Exception)
|
||||
capnp.AsyncIoContext moveAsyncContext"std::move"(capnp.AsyncIoContext)
|
||||
|
||||
|
||||
cdef extern from "<capnp/pretty-print.h>" namespace " ::capnp":
|
||||
StringTree printStructReader" ::capnp::prettyPrint"(C_DynamicStruct.Reader) except +reraise_kj_exception
|
||||
StringTree printStructBuilder" ::capnp::prettyPrint"(DynamicStruct_Builder) except +reraise_kj_exception
|
||||
@@ -1289,7 +1264,8 @@ cdef class _DynamicStructBuilder:
|
||||
:Raises: :exc:`KjException` if this isn't the message's root struct.
|
||||
"""
|
||||
self._check_write()
|
||||
await _VoidPromise()._init(writeMessage(deref(stream.thisptr.get()), deref((<_MessageBuilder>self._parent).thisptr)))
|
||||
await _voidpromise_to_asyncio(
|
||||
writeMessage(deref(stream.thisptr.get()), deref((<_MessageBuilder>self._parent).thisptr)))
|
||||
self._is_written = True
|
||||
|
||||
def write_packed(self, file):
|
||||
@@ -1657,12 +1633,12 @@ cdef class _DynamicStructPipeline:
|
||||
|
||||
cdef class _DynamicOrphan:
|
||||
cdef _init(self, C_DynamicOrphan other, object parent):
|
||||
self.thisptr = moveOrphan(other)
|
||||
self.thisptr = move(other)
|
||||
self._parent = parent
|
||||
return self
|
||||
|
||||
cdef C_DynamicOrphan move(self):
|
||||
return moveOrphan(self.thisptr)
|
||||
return move(self.thisptr)
|
||||
|
||||
cpdef get(self):
|
||||
"""Returns a python object corresponding to the DynamicValue owned by this orphan
|
||||
@@ -1831,11 +1807,8 @@ def _asyncio_close_patch(loop, oldclose, _EventLoop kjloop):
|
||||
|
||||
cdef class _EventLoop:
|
||||
cdef object __weakref__ # Needed to make this class weak-referenceable
|
||||
cdef Own[LowLevelAsyncIoProvider] lowLevelProvider
|
||||
cdef Own[AsyncIoProvider] provider
|
||||
cdef WaitScope * waitScope
|
||||
|
||||
cdef AsyncIoEventPort *customPort
|
||||
cdef WaitScope* waitScope
|
||||
cdef AsyncIoEventPort* customPort
|
||||
|
||||
def __init__(self):
|
||||
self._init()
|
||||
@@ -1848,10 +1821,8 @@ cdef class _EventLoop:
|
||||
loop.close = _partial(_asyncio_close_patch, loop, loop.close, self)
|
||||
|
||||
def __dealloc__(self):
|
||||
if not self.customPort == NULL:
|
||||
# If we have a custom port, the waitscope is not owned by provider, we have to delete it manually
|
||||
del self.waitScope
|
||||
del self.customPort
|
||||
del self.waitScope
|
||||
del self.customPort
|
||||
|
||||
|
||||
_C_DEFAULT_EVENT_LOOP_LOCAL = _threading.local()
|
||||
@@ -1882,7 +1853,7 @@ cdef class _CallContext:
|
||||
cdef CallContext * thisptr
|
||||
|
||||
cdef _init(self, CallContext other):
|
||||
self.thisptr = new CallContext(moveCallContext(other))
|
||||
self.thisptr = new CallContext(move(other))
|
||||
return self
|
||||
|
||||
def __dealloc__(self):
|
||||
@@ -1906,121 +1877,52 @@ cdef class _CallContext:
|
||||
self.thisptr.allowCancellation()
|
||||
|
||||
cpdef tail_call(self, _Request tailRequest):
|
||||
promise = _VoidPromise()._init(self.thisptr.tailCall(moveRequest(deref(tailRequest.thisptr_child))))
|
||||
return promise
|
||||
return _voidpromise_to_asyncio(self.thisptr.tailCall(move(deref(tailRequest.thisptr_child))))
|
||||
|
||||
cdef void _promise_check_consumed(PromiseTypes promise) except*:
|
||||
if promise.thisptr.get() == NULL:
|
||||
raise KjException(
|
||||
"Promise was already used in a consuming operation. You can no longer use this Promise object")
|
||||
|
||||
cdef _promise_then(PromiseTypes self, func, error_func, num_args, attach=None) except +reraise_kj_exception:
|
||||
_promise_check_consumed(self)
|
||||
|
||||
argspec = None
|
||||
try:
|
||||
argspec = _inspect.getfullargspec(func)
|
||||
except (TypeError, ValueError):
|
||||
pass
|
||||
if argspec:
|
||||
args_length = len(argspec.args) if argspec.args else 0
|
||||
defaults_length = len(argspec.defaults) if argspec.defaults else 0
|
||||
if args_length - defaults_length != num_args:
|
||||
raise KjException(f'Function passed to `then` call must take exactly {num_args} arguments')
|
||||
|
||||
return _Promise()._init(
|
||||
helpers.then(move(self.thisptr), capnp.heap[PyRefCounter](<PyObject *>func),
|
||||
capnp.heap[PyRefCounter](<PyObject *>error_func))
|
||||
.attach(capnp.heap[PyRefCounter](<PyObject *> attach)))
|
||||
|
||||
cdef _promise_to_asyncio(PromiseTypes promise):
|
||||
_promise_check_consumed(promise)
|
||||
|
||||
cdef _promise_to_asyncio(PyPromise promise):
|
||||
fut = asyncio.get_running_loop().create_future()
|
||||
def success(res): return fut.set_result(res) if not fut.cancelled() else None
|
||||
def exception(err): return fut.set_exception(err) if not fut.cancelled() else None
|
||||
def done(fut): return fut.kjpromise.cancel() if fut.cancelled() else None
|
||||
# Attach the promise to the future, so that it doesn't get destroyed
|
||||
fut.kjpromise = _promise_then(
|
||||
promise,
|
||||
lambda res: fut.set_result(res) if not fut.cancelled() else None,
|
||||
lambda err: fut.set_exception(err) if not fut.cancelled() else None,
|
||||
1)
|
||||
del promise
|
||||
fut.add_done_callback(
|
||||
lambda fut: fut.kjpromise.cancel() if fut.cancelled() else None)
|
||||
fut.kjpromise = _Promise()._init(helpers.then(
|
||||
move(promise),
|
||||
capnp.heap[PyRefCounter](<PyObject *>success),
|
||||
capnp.heap[PyRefCounter](<PyObject *>exception)))
|
||||
fut.add_done_callback(done)
|
||||
return fut
|
||||
|
||||
cdef _voidpromise_to_asyncio(VoidPromise promise):
|
||||
return _promise_to_asyncio(helpers.convert_to_pypromise(move(promise)))
|
||||
|
||||
cdef class _Promise:
|
||||
cdef Own[PyPromise] thisptr
|
||||
|
||||
def __init__(self, obj=None):
|
||||
C_DEFAULT_EVENT_LOOP_GETTER()
|
||||
if obj is not None:
|
||||
self.thisptr = capnp.heap[PyPromise](capnp.heap[PyRefCounter](<PyObject *>obj))
|
||||
|
||||
cdef _init(self, PyPromise other):
|
||||
self.thisptr = capnp.heap[PyPromise](movePromise(other))
|
||||
self.thisptr = capnp.heap[PyPromise](move(other))
|
||||
return self
|
||||
|
||||
async def a_wait(self):
|
||||
"""
|
||||
Asyncio version of wait().
|
||||
Required when using asyncio for socket communication.
|
||||
|
||||
Will still work with non-asyncio socket communication, but requires async handling of the function call.
|
||||
"""
|
||||
return await _promise_to_asyncio(self)
|
||||
|
||||
def __await__(self):
|
||||
return _promise_to_asyncio(self).__await__()
|
||||
|
||||
cpdef cancel(self) except +reraise_kj_exception:
|
||||
self.thisptr = Own[PyPromise]()
|
||||
|
||||
|
||||
cdef class _VoidPromise:
|
||||
cdef Own[VoidPromise] thisptr
|
||||
|
||||
|
||||
cdef _init(self, VoidPromise other):
|
||||
C_DEFAULT_EVENT_LOOP_GETTER()
|
||||
self.thisptr = capnp.heap[VoidPromise](moveVoidPromise(other))
|
||||
return self
|
||||
|
||||
async def a_wait(self):
|
||||
"""
|
||||
Asyncio version of wait().
|
||||
Required when using asyncio for socket communication.
|
||||
|
||||
Will still work with non-asyncio socket communication, but requires async handling of the function call.
|
||||
"""
|
||||
# TODO: Is keeping a separate _VoidPromise class really worth it? Does it make things faster?
|
||||
return await _promise_to_asyncio[_Promise](self.as_pypromise())
|
||||
|
||||
def __await__(self):
|
||||
return _promise_to_asyncio[_Promise](self.as_pypromise()).__await__()
|
||||
|
||||
cpdef as_pypromise(self) except +reraise_kj_exception:
|
||||
_promise_check_consumed(self)
|
||||
return _Promise()._init(helpers.convert_to_pypromise(move(self.thisptr)))
|
||||
|
||||
cpdef cancel(self) except +reraise_kj_exception:
|
||||
self.thisptr = Own[VoidPromise]()
|
||||
|
||||
|
||||
|
||||
cdef class _RemotePromise:
|
||||
cdef object _parent
|
||||
"""A pointer to a parent object that needs to be kept alive for this promise to function.
|
||||
Note that _Promise and _VoidPromise don't have such pointer. The reason is that in _RemotePromise
|
||||
the parent pointer needs to be passed around through _RemotePromise._get. If an object needs to
|
||||
be kept alive in _Promise or _VoidPromise, it can be attached to the underlying C++ promise."""
|
||||
"""A pointer to a parent object that needs to be kept alive for this promise to function."""
|
||||
|
||||
cdef Own[RemotePromise] thisptr
|
||||
|
||||
cdef _init(self, RemotePromise other, object parent=None):
|
||||
self.thisptr = capnp.heap[RemotePromise](moveRemotePromise(other))
|
||||
self.thisptr = capnp.heap[RemotePromise](move(other))
|
||||
self._parent = parent
|
||||
return self
|
||||
|
||||
cdef void _check_consumed(self) except*:
|
||||
if self.thisptr.get() == NULL:
|
||||
raise KjException(
|
||||
"Promise was already used in a consuming operation. You can no longer use this Promise object")
|
||||
|
||||
async def a_wait(self):
|
||||
"""
|
||||
Asyncio version of wait().
|
||||
@@ -2028,20 +1930,17 @@ cdef class _RemotePromise:
|
||||
|
||||
Will still work with non-asyncio socket communication, but requires async handling of the function call.
|
||||
"""
|
||||
return await _promise_to_asyncio(self)
|
||||
self._check_consumed()
|
||||
cdef Own[RemotePromise] thisptr = move(self.thisptr)
|
||||
return await _promise_to_asyncio(helpers.convert_to_pypromise(move(deref(thisptr))))
|
||||
|
||||
def __await__(self):
|
||||
return _promise_to_asyncio(self).__await__()
|
||||
|
||||
cpdef as_pypromise(self) except +reraise_kj_exception:
|
||||
_promise_check_consumed(self)
|
||||
parent = self._parent
|
||||
self._parent = None # We don't need parent anymore. Setting to none allows quicker garbage collection
|
||||
return _Promise()._init(helpers.convert_to_pypromise(move(self.thisptr))
|
||||
.attach(capnp.heap[PyRefCounter](<PyObject *>parent)))
|
||||
self._check_consumed()
|
||||
cdef Own[RemotePromise] thisptr = move(self.thisptr)
|
||||
return _promise_to_asyncio(helpers.convert_to_pypromise(move(deref(thisptr)))).__await__()
|
||||
|
||||
cpdef _get(self, field) except +reraise_kj_exception:
|
||||
_promise_check_consumed(self)
|
||||
self._check_consumed()
|
||||
cdef int type = (<C_DynamicValue.Pipeline>self.thisptr.get().get(field)).getType()
|
||||
if type == capnp.TYPE_CAPABILITY:
|
||||
return _DynamicCapabilityClient()._init(
|
||||
@@ -2064,7 +1963,7 @@ cdef class _RemotePromise:
|
||||
property schema:
|
||||
"""A property that returns the _StructSchema object matching this reader"""
|
||||
def __get__(self):
|
||||
_promise_check_consumed(self)
|
||||
self._check_consumed()
|
||||
return _StructSchema()._init_child(self.thisptr.get().getSchema())
|
||||
|
||||
def __dir__(self):
|
||||
@@ -2083,7 +1982,7 @@ cdef class _Request(_DynamicStructBuilder):
|
||||
cdef public bint is_consumed
|
||||
|
||||
cdef _init_child(self, Request other, parent):
|
||||
self.thisptr_child = new Request(moveRequest(other))
|
||||
self.thisptr_child = new Request(move(other))
|
||||
self._init(<DynamicStruct_Builder>deref(self.thisptr_child), parent)
|
||||
self.is_consumed = False
|
||||
return self
|
||||
@@ -2102,7 +2001,7 @@ cdef class _Response(_DynamicStructReader):
|
||||
cdef Response * thisptr_child
|
||||
|
||||
cdef _init_child(self, Response other, parent):
|
||||
self.thisptr_child = new Response(moveResponse(other))
|
||||
self.thisptr_child = new Response(move(other))
|
||||
self._init(<C_DynamicStruct.Reader>deref(self.thisptr_child), parent)
|
||||
return self
|
||||
|
||||
@@ -2276,11 +2175,11 @@ cdef class _TwoPartyVatNetwork:
|
||||
|
||||
cdef _init(self, _AsyncIoStream stream, Side side, schema_cpp.ReaderOptions opts):
|
||||
self.stream = stream
|
||||
self.thisptr = makeTwoPartyVatNetwork(deref(stream.thisptr), side, opts)
|
||||
self.thisptr = capnp.heap[C_TwoPartyVatNetwork](deref(stream.thisptr), side, opts)
|
||||
return self
|
||||
|
||||
cpdef on_disconnect(self) except +reraise_kj_exception:
|
||||
return _VoidPromise()._init(deref(self.thisptr).onDisconnect())
|
||||
return _voidpromise_to_asyncio(deref(self.thisptr).onDisconnect())
|
||||
|
||||
|
||||
cdef class TwoPartyClient:
|
||||
@@ -2291,7 +2190,7 @@ cdef class TwoPartyClient:
|
||||
:param traversal_limit_in_words: Pointer derefence limit (see https://capnproto.org/cxx.html).
|
||||
:param nesting_limit: Recursive limit when reading types (see https://capnproto.org/cxx.html).
|
||||
"""
|
||||
cdef RpcSystem * thisptr
|
||||
cdef Own[RpcSystem] thisptr
|
||||
cdef _TwoPartyVatNetwork _network
|
||||
|
||||
def __init__(self, socket=None, traversal_limit_in_words=None, nesting_limit=None):
|
||||
@@ -2303,17 +2202,13 @@ cdef class TwoPartyClient:
|
||||
else:
|
||||
raise ValueError(f"Argument socket should be a AsyncIoStream, was {type(socket)}")
|
||||
|
||||
self.thisptr = new RpcSystem(makeRpcClient(deref(self._network.thisptr)))
|
||||
|
||||
def __dealloc__(self):
|
||||
if not self.thisptr == NULL:
|
||||
del self.thisptr
|
||||
self.thisptr = capnp.heap[RpcSystem](makeRpcClient(deref(self._network.thisptr)))
|
||||
|
||||
cpdef bootstrap(self) except +reraise_kj_exception:
|
||||
return _CapabilityClient()._init(helpers.bootstrapHelper(deref(self.thisptr)), self)
|
||||
|
||||
cpdef on_disconnect(self) except +reraise_kj_exception:
|
||||
return _VoidPromise()._init(deref(self._network.thisptr).onDisconnect())
|
||||
return self._network.on_disconnect()
|
||||
|
||||
|
||||
cdef class TwoPartyServer:
|
||||
@@ -2325,7 +2220,7 @@ cdef class TwoPartyServer:
|
||||
:param traversal_limit_in_words: Pointer derefence limit (see https://capnproto.org/cxx.html).
|
||||
:param nesting_limit: Recursive limit when reading types (see https://capnproto.org/cxx.html).
|
||||
"""
|
||||
cdef RpcSystem * thisptr
|
||||
cdef Own[RpcSystem] thisptr
|
||||
cdef _TwoPartyVatNetwork _network
|
||||
|
||||
def __init__(self, socket=None, bootstrap=None, traversal_limit_in_words=None, nesting_limit=None):
|
||||
@@ -2339,18 +2234,15 @@ cdef class TwoPartyServer:
|
||||
raise ValueError(f"Argument socket should be a AsyncIoStream, was {type(socket)}")
|
||||
|
||||
cdef _InterfaceSchema schema = bootstrap.schema
|
||||
self.thisptr = new RpcSystem(makeRpcServerBootstrap(
|
||||
self.thisptr = capnp.heap[RpcSystem](makeRpcServerBootstrap(
|
||||
deref(self._network.thisptr), helpers.server_to_client(schema.thisptr, <PyObject *>bootstrap)))
|
||||
|
||||
def __dealloc__(self):
|
||||
del self.thisptr
|
||||
|
||||
cpdef on_disconnect(self) except +reraise_kj_exception:
|
||||
return _VoidPromise()._init(deref(self._network.thisptr).onDisconnect())
|
||||
|
||||
cpdef bootstrap(self) except +reraise_kj_exception:
|
||||
return _CapabilityClient()._init(helpers.bootstrapHelperServer(deref(self.thisptr)), self)
|
||||
|
||||
cpdef on_disconnect(self) except +reraise_kj_exception:
|
||||
return self._network.on_disconnect()
|
||||
|
||||
|
||||
cdef class _AsyncIoStream:
|
||||
cdef Own[AsyncIoStream] thisptr
|
||||
@@ -3119,7 +3011,7 @@ class _StructModule(object):
|
||||
|
||||
:rtype: :class:`_DynamicStructReader`"""
|
||||
cdef schema_cpp.ReaderOptions opts = make_reader_opts(traversal_limit_in_words, nesting_limit)
|
||||
reader = await _Promise()._init(tryReadMessage(deref(stream.thisptr.get()), opts))
|
||||
reader = await _promise_to_asyncio(tryReadMessage(deref(stream.thisptr.get()), opts))
|
||||
if reader is None:
|
||||
return
|
||||
return reader.get_root(self.schema)
|
||||
|
||||
Reference in New Issue
Block a user