Fix memleak and simplify dynamic client api
This commit is contained in:
@@ -12,8 +12,14 @@ class PythonInterfaceDynamicImpl final: public capnp::DynamicCapability::Server
|
|||||||
public:
|
public:
|
||||||
PyObject * py_server;
|
PyObject * py_server;
|
||||||
|
|
||||||
PythonInterfaceDynamicImpl(capnp::InterfaceSchema & schema, PyObject * py_server)
|
PythonInterfaceDynamicImpl(capnp::InterfaceSchema & schema, PyObject * _py_server)
|
||||||
: capnp::DynamicCapability::Server(schema), py_server(py_server) {}
|
: capnp::DynamicCapability::Server(schema), py_server(_py_server) {
|
||||||
|
Py_INCREF(_py_server);
|
||||||
|
}
|
||||||
|
|
||||||
|
~PythonInterfaceDynamicImpl() {
|
||||||
|
Py_DECREF(py_server);
|
||||||
|
}
|
||||||
|
|
||||||
kj::Promise<void> call(capnp::InterfaceSchema::Method method,
|
kj::Promise<void> call(capnp::InterfaceSchema::Method method,
|
||||||
capnp::CallContext< capnp::DynamicStruct, capnp::DynamicStruct> context) {
|
capnp::CallContext< capnp::DynamicStruct, capnp::DynamicStruct> context) {
|
||||||
|
|||||||
143
capnp/capnp.pyx
143
capnp/capnp.pyx
@@ -9,7 +9,7 @@
|
|||||||
cimport cython
|
cimport cython
|
||||||
cimport capnp_cpp as capnp
|
cimport capnp_cpp as capnp
|
||||||
cimport schema_cpp
|
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, 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, CallContext
|
||||||
|
|
||||||
from schema_cpp cimport Node as C_Node, EnumNode as C_EnumNode
|
from schema_cpp cimport Node as C_Node, EnumNode as C_EnumNode
|
||||||
from cython.operator cimport dereference as deref
|
from cython.operator cimport dereference as deref
|
||||||
@@ -34,6 +34,12 @@ ctypedef double Float64
|
|||||||
from libc.stdlib cimport malloc, free
|
from libc.stdlib cimport malloc, free
|
||||||
from libcpp cimport bool as cbool
|
from libcpp cimport bool as cbool
|
||||||
|
|
||||||
|
from types import ModuleType as _ModuleType
|
||||||
|
import os as _os
|
||||||
|
import sys as _sys
|
||||||
|
import imp as _imp
|
||||||
|
from functools import partial as _partial
|
||||||
|
|
||||||
# By making it public, we'll be able to call it from capabilityHelper.h
|
# By making it public, we'll be able to call it from capabilityHelper.h
|
||||||
cdef public object wrap_dynamic_struct_reader(C_DynamicStruct.Reader & reader):
|
cdef public object wrap_dynamic_struct_reader(C_DynamicStruct.Reader & reader):
|
||||||
return _DynamicStructReader()._init(reader, None)
|
return _DynamicStructReader()._init(reader, None)
|
||||||
@@ -95,6 +101,7 @@ cdef extern from "capnp/list.h" namespace " ::capnp":
|
|||||||
cdef extern from "<utility>" namespace "std":
|
cdef extern from "<utility>" namespace "std":
|
||||||
C_DynamicOrphan moveOrphan"std::move"(C_DynamicOrphan)
|
C_DynamicOrphan moveOrphan"std::move"(C_DynamicOrphan)
|
||||||
Request moveRequest"std::move"(Request)
|
Request moveRequest"std::move"(Request)
|
||||||
|
Response moveResponse"std::move"(Response)
|
||||||
PyPromise movePromise"std::move"(PyPromise)
|
PyPromise movePromise"std::move"(PyPromise)
|
||||||
RemotePromise moveRemotePromise"std::move"(RemotePromise)
|
RemotePromise moveRemotePromise"std::move"(RemotePromise)
|
||||||
CallContext moveCallContext"std::move"(CallContext)
|
CallContext moveCallContext"std::move"(CallContext)
|
||||||
@@ -969,7 +976,7 @@ cdef class EventLoop:
|
|||||||
if promise.is_consumed:
|
if promise.is_consumed:
|
||||||
raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
raise RuntimeError('Promise was already used in a consuming operation. You can no longer use this Promise object')
|
||||||
|
|
||||||
ret = _DynamicStructReader()._init(self.thisptr.wait_remote(moveRemotePromise(deref(promise.thisptr))), promise._parent)
|
ret = _Response()._init_child(self.thisptr.wait_remote(moveRemotePromise(deref(promise.thisptr))), promise._parent)
|
||||||
promise.is_consumed = True
|
promise.is_consumed = True
|
||||||
|
|
||||||
return ret
|
return ret
|
||||||
@@ -982,108 +989,24 @@ cdef class EventLoop:
|
|||||||
# Py_INCREF(error_func)
|
# Py_INCREF(error_func)
|
||||||
# return Promise()._init(capnp.there(self.thisptr, deref(promise.thisptr), <PyObject *>func, <PyObject *>error_func))
|
# return Promise()._init(capnp.there(self.thisptr, deref(promise.thisptr), <PyObject *>func, <PyObject *>error_func))
|
||||||
|
|
||||||
cdef class _Request:
|
cdef class _Request(_DynamicStructBuilder):
|
||||||
cdef Request * thisptr
|
cdef Request * thisptr_child
|
||||||
cdef public object _parent
|
|
||||||
|
|
||||||
cdef _init(self, Request other, parent):
|
cdef _init_child(self, Request other, parent):
|
||||||
self.thisptr = new Request(moveRequest(other))
|
self.thisptr_child = new Request(moveRequest(other))
|
||||||
self._parent = parent
|
self._init(<C_DynamicStruct.Builder>deref(self.thisptr_child), parent)
|
||||||
return self
|
return self
|
||||||
|
|
||||||
cpdef send(self):
|
cpdef send(self):
|
||||||
return _RemotePromise()._init(self.thisptr.send(), self._parent)
|
return _RemotePromise()._init(self.thisptr_child.send(), self._parent)
|
||||||
cdef _get(self, field):
|
|
||||||
cdef C_DynamicValue.Builder value = self.thisptr.get(field)
|
|
||||||
|
|
||||||
return to_python_builder(value, self._parent)
|
cdef class _Response(_DynamicStructReader):
|
||||||
|
cdef Response * thisptr_child
|
||||||
|
|
||||||
def __getattr__(self, field):
|
cdef _init_child(self, Response other, parent):
|
||||||
return self._get(field)
|
self.thisptr_child = new Response(moveResponse(other))
|
||||||
|
self._init(<C_DynamicStruct.Reader>deref(self.thisptr_child), parent)
|
||||||
def __setattr__(self, field, value):
|
return self
|
||||||
_setDynamicFieldPtr(self.thisptr, field, value, self._parent)
|
|
||||||
|
|
||||||
def _has(self, field):
|
|
||||||
return self.thisptr.has(field)
|
|
||||||
|
|
||||||
cpdef init(self, field, size=None):
|
|
||||||
"""Method for initializing fields that are of type union/struct/list
|
|
||||||
|
|
||||||
Typically, you don't have to worry about initializing structs/unions, so this method is mainly for lists.
|
|
||||||
|
|
||||||
:type field: str
|
|
||||||
:param field: The field name to initialize
|
|
||||||
|
|
||||||
:type size: int
|
|
||||||
:param size: The size of the list to initiialize. This should be None for struct/union initialization.
|
|
||||||
|
|
||||||
:rtype: :class:`_DynamicStructBuilder` or :class:`_DynamicListBuilder`
|
|
||||||
|
|
||||||
:Raises: :exc:`exceptions.ValueError` if the field isn't in this struct
|
|
||||||
"""
|
|
||||||
if size is None:
|
|
||||||
return to_python_builder(self.thisptr.init(field), self._parent)
|
|
||||||
else:
|
|
||||||
return to_python_builder(self.thisptr.init(field, size), self._parent)
|
|
||||||
|
|
||||||
cpdef init_resizable_list(self, field):
|
|
||||||
"""Method for initializing fields that are of type list (of structs)
|
|
||||||
|
|
||||||
This version of init returns a :class:`_DynamicResizableListBuilder` that allows you to add members one at a time (ie. if you don't know the size for sure). This is only meant for lists of Cap'n Proto objects, since for primitive types you can just define a normal python list and fill it yourself.
|
|
||||||
|
|
||||||
.. warning:: You need to call :meth:`_DynamicResizableListBuilder.finish` on the list object before serializing the Cap'n Proto message. Failure to do so will cause your objects not to be written out as well as leaking orphan structs into your message.
|
|
||||||
|
|
||||||
:type field: str
|
|
||||||
:param field: The field name to initialize
|
|
||||||
|
|
||||||
:rtype: :class:`_DynamicResizableListBuilder`
|
|
||||||
|
|
||||||
:Raises: :exc:`exceptions.ValueError` if the field isn't in this struct
|
|
||||||
"""
|
|
||||||
return _DynamicResizableListBuilder(self, field, _StructSchema()._init((<C_DynamicValue.Builder>self.thisptr.get(field)).asList().getStructElementType()))
|
|
||||||
|
|
||||||
cpdef which(self):
|
|
||||||
"""Returns the enum corresponding to the union in this struct
|
|
||||||
|
|
||||||
Enums are just strings in the python Cap'n Proto API, so this function will either return a string equal to the field name of the active field in the union, or throw a ValueError if this isn't a union, or a struct with an unnamed union::
|
|
||||||
|
|
||||||
person = addressbook.Person.new_message()
|
|
||||||
|
|
||||||
person.which()
|
|
||||||
# ValueError: member was null
|
|
||||||
|
|
||||||
a.employment.employer = 'foo'
|
|
||||||
print employment.which()
|
|
||||||
# 'employer'
|
|
||||||
|
|
||||||
:rtype: str
|
|
||||||
:return: A string/enum corresponding to what field is set in the union
|
|
||||||
|
|
||||||
:Raises: :exc:`exceptions.ValueError` if this struct doesn't contain a union
|
|
||||||
"""
|
|
||||||
cdef object which = getEnumString(deref(self.thisptr))
|
|
||||||
if len(which) == 0:
|
|
||||||
raise ValueError("Attempted to call which on a non-union type")
|
|
||||||
|
|
||||||
return which
|
|
||||||
|
|
||||||
property schema:
|
|
||||||
"""A property that returns the _StructSchema object matching this writer"""
|
|
||||||
def __get__(self):
|
|
||||||
return _StructSchema()._init(self.thisptr.getSchema())
|
|
||||||
|
|
||||||
def __dir__(self):
|
|
||||||
return list(self.schema.fieldnames)
|
|
||||||
|
|
||||||
def __str__(self):
|
|
||||||
return printRequest(deref(self.thisptr)).flatten().cStr()
|
|
||||||
|
|
||||||
def __repr__(self):
|
|
||||||
return '<%s builder %s>' % (self.schema.node.displayName, strRequest(deref(self.thisptr)).cStr())
|
|
||||||
|
|
||||||
def to_dict(self):
|
|
||||||
return _to_dict(self)
|
|
||||||
|
|
||||||
cdef class _DynamicCapabilityClient:
|
cdef class _DynamicCapabilityClient:
|
||||||
cdef C_DynamicCapability.Client thisptr
|
cdef C_DynamicCapability.Client thisptr
|
||||||
@@ -1101,7 +1024,7 @@ cdef class _DynamicCapabilityClient:
|
|||||||
self.thisptr = new_client(s.thisptr, <PyObject *>server, loop.thisptr)
|
self.thisptr = new_client(s.thisptr, <PyObject *>server, loop.thisptr)
|
||||||
self._server = server
|
self._server = server
|
||||||
|
|
||||||
cpdef _new_request_helper(self, name, firstSegmentWordSize, kwargs) except +ValueError:
|
cpdef _send_helper(self, name, firstSegmentWordSize, kwargs) except +ValueError:
|
||||||
cdef Request * request = new Request(self.thisptr.newRequest(name, firstSegmentWordSize))
|
cdef Request * request = new Request(self.thisptr.newRequest(name, firstSegmentWordSize))
|
||||||
|
|
||||||
for key, val in kwargs.items():
|
for key, val in kwargs.items():
|
||||||
@@ -1109,11 +1032,20 @@ cdef class _DynamicCapabilityClient:
|
|||||||
|
|
||||||
return _RemotePromise()._init(request.send(), self)
|
return _RemotePromise()._init(request.send(), self)
|
||||||
|
|
||||||
cpdef request(self, name, firstSegmentWordSize=0) except +ValueError:
|
cpdef _request_helper(self, name, firstSegmentWordSize=0) except +ValueError:
|
||||||
return _Request()._init(self.thisptr.newRequest(name, firstSegmentWordSize), self)
|
return _Request()._init_child(self.thisptr.newRequest(name, firstSegmentWordSize), self)
|
||||||
|
|
||||||
def send(self, name, firstSegmentWordSize=0, **kwargs):
|
def _request(self, name, firstSegmentWordSize=0):
|
||||||
return self._new_request_helper(name, firstSegmentWordSize, kwargs)
|
return self._request_helper(name, firstSegmentWordSize)
|
||||||
|
|
||||||
|
def _send(self, name, *args, firstSegmentWordSize=0, **kwargs):
|
||||||
|
return self._send_helper(name, firstSegmentWordSize, kwargs)
|
||||||
|
|
||||||
|
def __getattr__(self, name):
|
||||||
|
if name.endswith('_request'):
|
||||||
|
short_name = name[:-8]
|
||||||
|
return _partial(self._request, short_name)
|
||||||
|
return _partial(self._send, name)
|
||||||
|
|
||||||
cdef class _Schema:
|
cdef class _Schema:
|
||||||
cdef C_Schema thisptr
|
cdef C_Schema thisptr
|
||||||
@@ -1630,11 +1562,6 @@ def _write_packed_message_to_fd(int fd, _MessageBuilder message):
|
|||||||
"""
|
"""
|
||||||
schema_cpp.writePackedMessageToFd(fd, deref(message.thisptr))
|
schema_cpp.writePackedMessageToFd(fd, deref(message.thisptr))
|
||||||
|
|
||||||
from types import ModuleType as _ModuleType
|
|
||||||
import os as _os
|
|
||||||
import sys as _sys
|
|
||||||
import imp as _imp
|
|
||||||
|
|
||||||
_global_schema_parser = None
|
_global_schema_parser = None
|
||||||
|
|
||||||
def load(file_name, display_name=None, imports=[]):
|
def load(file_name, display_name=None, imports=[]):
|
||||||
|
|||||||
@@ -150,7 +150,7 @@ cdef extern from "capnp/dynamic.h" namespace " ::capnp":
|
|||||||
|
|
||||||
cdef extern from "capnp/capability.h" namespace " ::capnp":
|
cdef extern from "capnp/capability.h" namespace " ::capnp":
|
||||||
cdef cppclass Response" ::capnp::Response< ::capnp::DynamicStruct>"(DynamicStruct.Reader):
|
cdef cppclass Response" ::capnp::Response< ::capnp::DynamicStruct>"(DynamicStruct.Reader):
|
||||||
pass
|
Response(Response)
|
||||||
cdef cppclass RemotePromise" ::capnp::RemotePromise< ::capnp::DynamicStruct>"(Promise[Response]):
|
cdef cppclass RemotePromise" ::capnp::RemotePromise< ::capnp::DynamicStruct>"(Promise[Response]):
|
||||||
RemotePromise(RemotePromise)
|
RemotePromise(RemotePromise)
|
||||||
|
|
||||||
@@ -291,7 +291,7 @@ cdef extern from "kj/async.h" namespace " ::kj":
|
|||||||
EventLoop()
|
EventLoop()
|
||||||
# Promise[void] yield_end'yield'()
|
# Promise[void] yield_end'yield'()
|
||||||
object wait(PyPromise) except+
|
object wait(PyPromise) except+
|
||||||
DynamicStruct.Reader wait_remote'wait'(RemotePromise) except+
|
Response wait_remote'wait'(RemotePromise)
|
||||||
object there(PyPromise) except+
|
object there(PyPromise) except+
|
||||||
PyPromise evalLater(PyObject * func)
|
PyPromise evalLater(PyObject * func)
|
||||||
PyPromise there(PyPromise, PyObject * func)
|
PyPromise there(PyPromise, PyObject * func)
|
||||||
|
|||||||
@@ -13,7 +13,6 @@ def example_client():
|
|||||||
client = example_capability_capnp.TestInterface.new_client(Server(), loop)
|
client = example_capability_capnp.TestInterface.new_client(Server(), loop)
|
||||||
|
|
||||||
req = client.request('foo')
|
req = client.request('foo')
|
||||||
req = client.request('foo2')
|
|
||||||
req.i = 5
|
req.i = 5
|
||||||
|
|
||||||
remote = req.send()
|
remote = req.send()
|
||||||
|
|||||||
@@ -17,14 +17,56 @@ def test_basic_client(capability):
|
|||||||
|
|
||||||
client = capability.TestInterface.new_client(Server(), loop)
|
client = capability.TestInterface.new_client(Server(), loop)
|
||||||
|
|
||||||
req = client.request('foo')
|
req = client._request('foo')
|
||||||
req.i = 5
|
req.i = 5
|
||||||
|
|
||||||
remote = req.send()
|
remote = req.send()
|
||||||
remote = client.send('foo', i=10)
|
|
||||||
response = loop.wait_remote(remote)
|
response = loop.wait_remote(remote)
|
||||||
|
|
||||||
# assert response.x == '26'
|
assert response.x == '26'
|
||||||
|
|
||||||
|
req = client.foo_request()
|
||||||
|
req.i = 5
|
||||||
|
|
||||||
|
remote = req.send()
|
||||||
|
response = loop.wait_remote(remote)
|
||||||
|
|
||||||
|
assert response.x == '26'
|
||||||
|
|
||||||
with pytest.raises(ValueError):
|
with pytest.raises(ValueError):
|
||||||
client.request('foo2')
|
client.foo2_request()
|
||||||
|
|
||||||
|
req = client.foo_request()
|
||||||
|
|
||||||
|
with pytest.raises(ValueError):
|
||||||
|
req.i = 'foo'
|
||||||
|
|
||||||
|
req = client.foo_request()
|
||||||
|
|
||||||
|
with pytest.raises(ValueError):
|
||||||
|
req.baz = 1
|
||||||
|
|
||||||
|
def test_simple_client(capability):
|
||||||
|
loop = capnp.EventLoop()
|
||||||
|
|
||||||
|
client = capability.TestInterface.new_client(Server(), loop)
|
||||||
|
|
||||||
|
remote = client._send('foo', i=5)
|
||||||
|
response = loop.wait_remote(remote)
|
||||||
|
|
||||||
|
assert response.x == '26'
|
||||||
|
|
||||||
|
|
||||||
|
remote = client.foo(i=5)
|
||||||
|
response = loop.wait_remote(remote)
|
||||||
|
|
||||||
|
assert response.x == '26'
|
||||||
|
|
||||||
|
with pytest.raises(ValueError):
|
||||||
|
remote = client.foo(i='foo')
|
||||||
|
|
||||||
|
with pytest.raises(ValueError):
|
||||||
|
remote = client.foo2(i=5)
|
||||||
|
|
||||||
|
with pytest.raises(ValueError):
|
||||||
|
remote = client.foo(baz=5)
|
||||||
|
|||||||
Reference in New Issue
Block a user