Pipelining almost completely wrapped
Waiting on some upstream changes in C++ libcapnp before I can finish
This commit is contained in:
112
capnp/capnp.pyx
112
capnp/capnp.pyx
@@ -9,7 +9,7 @@
|
||||
cimport cython
|
||||
cimport capnp_cpp as capnp
|
||||
cimport schema_cpp
|
||||
from capnp_cpp cimport Schema as C_Schema, StructSchema as C_StructSchema, InterfaceSchema as C_InterfaceSchema, DynamicStruct as C_DynamicStruct, DynamicValue as C_DynamicValue, Type as C_Type, DynamicList as C_DynamicList, fixMaybe, getEnumString, SchemaParser as C_SchemaParser, ParsedSchema as C_ParsedSchema, VOID, ArrayPtr, StringPtr, String, StringTree, DynamicOrphan as C_DynamicOrphan, ObjectPointer as C_DynamicObject, WordArrayPtr, DynamicCapability as C_DynamicCapability, new_client, Request, Response, RemotePromise, convert_to_pypromise, SimpleEventLoop, PyPromise, CallContext
|
||||
from capnp_cpp cimport Schema as C_Schema, StructSchema as C_StructSchema, InterfaceSchema as C_InterfaceSchema, DynamicStruct as C_DynamicStruct, DynamicValue as C_DynamicValue, Type as C_Type, DynamicList as C_DynamicList, fixMaybe, getEnumString, SchemaParser as C_SchemaParser, ParsedSchema as C_ParsedSchema, VOID, ArrayPtr, StringPtr, String, StringTree, DynamicOrphan as C_DynamicOrphan, ObjectPointer as C_DynamicObject, WordArrayPtr, DynamicCapability as C_DynamicCapability, new_client, Request, Response, RemotePromise, convert_to_pypromise, SimpleEventLoop, PyPromise, VoidPromise, CallContext
|
||||
|
||||
from schema_cpp cimport Node as C_Node, EnumNode as C_EnumNode
|
||||
from cython.operator cimport dereference as deref
|
||||
@@ -44,14 +44,28 @@ from functools import partial as _partial
|
||||
cdef public object wrap_dynamic_struct_reader(C_DynamicStruct.Reader & reader):
|
||||
return _DynamicStructReader()._init(reader, None)
|
||||
|
||||
cdef public void call_server_method(PyObject * _server, char * _method_name, CallContext & _context):
|
||||
cdef public void wrap_remote_call(PyObject * func, Response & r):
|
||||
response = _Response()._init_childptr(new Response(moveResponse(r)), None)
|
||||
|
||||
func_obj = <object>func
|
||||
# TODO: decref func?
|
||||
func_obj(response)
|
||||
|
||||
cdef public VoidPromise * call_server_method(PyObject * _server, char * _method_name, CallContext & _context):
|
||||
server = <object>_server
|
||||
method_name = <object>_method_name
|
||||
|
||||
context = _CallContext()._init(_context)
|
||||
getattr(server, method_name)(context)
|
||||
ret = getattr(server, method_name)(context)
|
||||
|
||||
# By making it public, we'll be able to call it from asyncHelper.h
|
||||
if ret is not None:
|
||||
if type(ret) is _VoidPromise:
|
||||
return new VoidPromise(moveVoidPromise(deref((<_VoidPromise>ret).thisptr)))
|
||||
else:
|
||||
raise ValueError('Server function returned a value that was not a VoidPromise: ' + str(ret))
|
||||
|
||||
return NULL
|
||||
|
||||
cdef public object wrap_kj_exception(capnp.Exception & exception):
|
||||
return None # TODO
|
||||
|
||||
@@ -103,6 +117,7 @@ cdef extern from "<utility>" namespace "std":
|
||||
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)
|
||||
|
||||
@@ -301,8 +316,6 @@ cdef class _DynamicListBuilder:
|
||||
return self._get(index)
|
||||
|
||||
def __setitem__(self, index, value):
|
||||
# TODO: share code with _DynamicStructBuilder.__setattr__
|
||||
|
||||
size = self.thisptr.size()
|
||||
if index >= size:
|
||||
raise IndexError('Out of bounds')
|
||||
@@ -387,6 +400,8 @@ cdef to_python_reader(C_DynamicValue.Reader self, object parent):
|
||||
return None
|
||||
elif type == capnp.TYPE_OBJECT:
|
||||
return _DynamicObjectReader()._init(self.asObject(), parent)
|
||||
elif type == capnp.TYPE_CAPABILITY:
|
||||
return _DynamicCapabilityClient()._init(self.asCapability(), parent)
|
||||
elif type == capnp.TYPE_UNKNOWN:
|
||||
raise ValueError("Cannot convert type to Python. Type is unknown by capnproto library")
|
||||
else:
|
||||
@@ -417,6 +432,8 @@ cdef to_python_builder(C_DynamicValue.Builder self, object parent):
|
||||
return None
|
||||
elif type == capnp.TYPE_OBJECT:
|
||||
return _DynamicObjectBuilder()._init(self.asObject(), parent)
|
||||
elif type == capnp.TYPE_CAPABILITY:
|
||||
return _DynamicCapabilityClient()._init(self.asCapability(), parent)
|
||||
elif type == capnp.TYPE_UNKNOWN:
|
||||
raise ValueError("Cannot convert type to Python. Type is unknown by capnproto library")
|
||||
else:
|
||||
@@ -428,6 +445,9 @@ cdef C_DynamicValue.Reader _extract_dynamic_struct_builder(_DynamicStructBuilder
|
||||
cdef C_DynamicValue.Reader _extract_dynamic_struct_reader(_DynamicStructReader value):
|
||||
return C_DynamicValue.Reader(value.thisptr)
|
||||
|
||||
cdef C_DynamicValue.Reader _extract_dynamic_client(_DynamicCapabilityClient value):
|
||||
return C_DynamicValue.Reader(value.thisptr)
|
||||
|
||||
cdef _setDynamicField(_DynamicSetterClasses thisptr, field, value, parent):
|
||||
cdef C_DynamicValue.Reader temp
|
||||
value_type = type(value)
|
||||
@@ -458,6 +478,8 @@ cdef _setDynamicField(_DynamicSetterClasses thisptr, field, value, parent):
|
||||
thisptr.set(field, _extract_dynamic_struct_builder(value))
|
||||
elif value_type is _DynamicStructReader:
|
||||
thisptr.set(field, _extract_dynamic_struct_reader(value))
|
||||
elif value_type is _DynamicCapabilityClient:
|
||||
thisptr.set(field, _extract_dynamic_client(value))
|
||||
else:
|
||||
raise ValueError("Non primitive type")
|
||||
|
||||
@@ -491,6 +513,8 @@ cdef _setDynamicFieldPtr(_DynamicSetterClasses * thisptr, field, value, parent):
|
||||
thisptr.set(field, _extract_dynamic_struct_builder(value))
|
||||
elif value_type is _DynamicStructReader:
|
||||
thisptr.set(field, _extract_dynamic_struct_reader(value))
|
||||
elif value_type is _DynamicCapabilityClient:
|
||||
thisptr.set(field, _extract_dynamic_client(value))
|
||||
else:
|
||||
raise ValueError("Non primitive type")
|
||||
|
||||
@@ -898,6 +922,7 @@ cdef class _CallContext:
|
||||
|
||||
cdef class Promise:
|
||||
cdef PyPromise * thisptr
|
||||
cdef public bint is_consumed
|
||||
|
||||
def __init__(self):
|
||||
self.is_consumed = True
|
||||
@@ -928,6 +953,38 @@ cdef class Promise:
|
||||
|
||||
return Promise()._init(capnp.then(deref(self.thisptr), <PyObject *>func, <PyObject *>error_func))
|
||||
|
||||
cdef class _VoidPromise:
|
||||
cdef VoidPromise * thisptr
|
||||
cdef public bint is_consumed
|
||||
|
||||
def __init__(self):
|
||||
self.is_consumed = True
|
||||
|
||||
cdef _init(self, VoidPromise other):
|
||||
self.is_consumed = False
|
||||
self.thisptr = new VoidPromise(moveVoidPromise(other))
|
||||
return self
|
||||
|
||||
def __dealloc__(self):
|
||||
del self.thisptr
|
||||
|
||||
cpdef wait(self) except+:
|
||||
if self.is_consumed:
|
||||
raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
||||
|
||||
self.thisptr.wait()
|
||||
self.is_consumed = True
|
||||
|
||||
|
||||
# cpdef then(self, func, error_func=None) except+:
|
||||
# if self.is_consumed:
|
||||
# raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
||||
|
||||
# Py_INCREF(func)
|
||||
# Py_INCREF(error_func)
|
||||
|
||||
# return Promise()._init(capnp.then(deref(self.thisptr), <PyObject *>func, <PyObject *>error_func))
|
||||
|
||||
cdef class _RemotePromise:
|
||||
cdef RemotePromise * thisptr
|
||||
cdef public bint is_consumed
|
||||
@@ -957,14 +1014,28 @@ cdef class _RemotePromise:
|
||||
cpdef as_pypromise(self) except +:
|
||||
Promise()._init(convert_to_pypromise(deref(self.thisptr)))
|
||||
|
||||
# cpdef then(self, func, error_func=None) except+:
|
||||
# if self.is_consumed:
|
||||
# raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
||||
cpdef then(self, func, error_func=None) except+:
|
||||
if self.is_consumed:
|
||||
raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
||||
|
||||
# Py_INCREF(func)
|
||||
# Py_INCREF(error_func)
|
||||
Py_INCREF(func)
|
||||
Py_INCREF(error_func)
|
||||
|
||||
# return _RemotePromise()._init(capnp.then(deref(self.thisptr), <PyObject *>func, <PyObject *>error_func))
|
||||
return _VoidPromise()._init(capnp.then(deref(self.thisptr), <PyObject *>func, <PyObject *>error_func))
|
||||
|
||||
cpdef _get(self, field) except +ValueError:
|
||||
return _DynamicCapabilityClient()._init((<C_DynamicValue.Pipeline>self.thisptr.get(field)).asCapability(), self._parent)
|
||||
|
||||
def __getattr__(self, field):
|
||||
return self._get(field)
|
||||
|
||||
property schema:
|
||||
"""A property that returns the _StructSchema object matching this reader"""
|
||||
def __get__(self):
|
||||
return _StructSchema()._init(self.thisptr.getSchema())
|
||||
|
||||
def __dir__(self):
|
||||
return list(self.schema.fieldnames)
|
||||
|
||||
cdef class EventLoop:
|
||||
cdef SimpleEventLoop thisptr
|
||||
@@ -1008,11 +1079,21 @@ cdef class _Response(_DynamicStructReader):
|
||||
self._init(<C_DynamicStruct.Reader>deref(self.thisptr_child), parent)
|
||||
return self
|
||||
|
||||
cdef _init_childptr(self, Response * other, parent):
|
||||
self.thisptr_child = other
|
||||
self._init(<C_DynamicStruct.Reader>deref(self.thisptr_child), parent)
|
||||
return self
|
||||
|
||||
cdef class _DynamicCapabilityClient:
|
||||
cdef C_DynamicCapability.Client thisptr
|
||||
cdef public object _event_loop, _server
|
||||
cdef public object _event_loop, _server, _parent
|
||||
|
||||
def __init__(self, schema, server, event_loop):
|
||||
cdef _init(self, C_DynamicCapability.Client other, object parent):
|
||||
self.thisptr = other
|
||||
self._parent = parent
|
||||
return self
|
||||
|
||||
cdef _init_vals(self, schema, server, event_loop):
|
||||
cdef _InterfaceSchema s
|
||||
if hasattr(schema, 'schema'):
|
||||
s = schema.schema
|
||||
@@ -1023,6 +1104,7 @@ cdef class _DynamicCapabilityClient:
|
||||
self._event_loop = event_loop
|
||||
self.thisptr = new_client(s.thisptr, <PyObject *>server, loop.thisptr)
|
||||
self._server = server
|
||||
return self
|
||||
|
||||
cpdef _send_helper(self, name, firstSegmentWordSize, kwargs) except +ValueError:
|
||||
cdef Request * request = new Request(self.thisptr.newRequest(name, firstSegmentWordSize))
|
||||
@@ -1287,7 +1369,7 @@ cdef class SchemaParser:
|
||||
elif proto.isInterface:
|
||||
def new_client(bound_local_module):
|
||||
def helper(server, loop):
|
||||
return _DynamicCapabilityClient(bound_local_module, server, loop)
|
||||
return _DynamicCapabilityClient()._init_vals(bound_local_module, server, loop)
|
||||
return helper
|
||||
local_module.schema = schema.as_interface()
|
||||
local_module.new_client = new_client(local_module)
|
||||
|
||||
Reference in New Issue
Block a user