From 09705b6d1f85a10b22ecf9dc0d0788323f65ab8b Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Thu, 23 Jun 2016 14:15:31 -0700 Subject: [PATCH 1/9] Add support for segment (de)serialization. --- capnp/includes/schema_cpp.pxd | 16 +++++++++++ capnp/lib/capnp.pxd | 1 + capnp/lib/capnp.pyx | 51 +++++++++++++++++++++++++++++++++++ 3 files changed, 68 insertions(+) diff --git a/capnp/includes/schema_cpp.pxd b/capnp/includes/schema_cpp.pxd index 3716648..d9b691c 100644 --- a/capnp/includes/schema_cpp.pxd +++ b/capnp/includes/schema_cpp.pxd @@ -667,6 +667,8 @@ cdef extern from "capnp/message.h" namespace " ::capnp": DynamicStruct_Builder initRootDynamicStruct'initRoot< ::capnp::DynamicStruct>'(StructSchema) void setRootDynamicStruct'setRoot< ::capnp::DynamicStruct::Reader>'(DynamicStruct.Reader) + ConstWordArrayArrayPtr getSegmentsForOutput'getSegmentsForOutput'() + AnyPointer.Builder getRootAnyPointer'getRoot< ::capnp::AnyPointer>'() DynamicOrphan newOrphan'getOrphanage().newOrphan'(StructSchema) @@ -691,6 +693,10 @@ cdef extern from "capnp/message.h" namespace " ::capnp": MallocMessageBuilder() MallocMessageBuilder(int) + cdef cppclass SegmentArrayMessageReader(MessageReader): + SegmentArrayMessageReader(ConstWordArrayArrayPtr array) except +reraise_kj_exception + SegmentArrayMessageReader(ConstWordArrayArrayPtr array, ReaderOptions) except +reraise_kj_exception + cdef cppclass FlatMessageBuilder(MessageBuilder): FlatMessageBuilder(WordArrayPtr array) FlatMessageBuilder(WordArrayPtr array, ReaderOptions) @@ -714,6 +720,16 @@ cdef extern from "kj/common.h" namespace " ::kj": ByteArrayPtr(byte *, size_t size) size_t size() byte& operator[](size_t index) + cdef cppclass ConstWordArrayPtr " ::kj::ArrayPtr< const ::capnp::word>": + ConstWordArrayPtr() + ConstWordArrayPtr(word *, size_t size) + size_t size() + const word* begin() + cdef cppclass ConstWordArrayArrayPtr " ::kj::ArrayPtr< const ::kj::ArrayPtr< const ::capnp::word>>": + ConstWordArrayArrayPtr() + ConstWordArrayArrayPtr(ConstWordArrayPtr*, size_t size) + size_t size() + ConstWordArrayPtr& operator[](size_t index) cdef extern from "kj/array.h" namespace " ::kj": # Cython can't handle Array[word] as a function argument diff --git a/capnp/lib/capnp.pxd b/capnp/lib/capnp.pxd index 1673365..6ac296a 100644 --- a/capnp/lib/capnp.pxd +++ b/capnp/lib/capnp.pxd @@ -53,6 +53,7 @@ cdef class _DynamicStructBuilder: cdef _check_write(self) cpdef to_bytes(_DynamicStructBuilder self) except +reraise_kj_exception + cpdef to_segments(_DynamicStructBuilder self) except +reraise_kj_exception cpdef _to_bytes_packed_helper(_DynamicStructBuilder self, word_count) except +reraise_kj_exception cpdef to_bytes_packed(_DynamicStructBuilder self) except +reraise_kj_exception diff --git a/capnp/lib/capnp.pyx b/capnp/lib/capnp.pyx index a052146..8ee1c3d 100644 --- a/capnp/lib/capnp.pyx +++ b/capnp/lib/capnp.pyx @@ -1192,6 +1192,12 @@ cdef class _DynamicStructBuilder: self._is_written = True return ret + cpdef to_segments(_DynamicStructBuilder self) except +reraise_kj_exception: + self._check_write() + cdef _MessageBuilder builder = self._parent + segments = builder.get_segments_for_output() + return segments + cpdef _to_bytes_packed_helper(_DynamicStructBuilder self, word_count) except +reraise_kj_exception: cdef _MessageBuilder builder = self._parent array = helpers.messageToPackedBytes(deref(builder.thisptr), word_count) @@ -2989,6 +2995,9 @@ class _StructModule(object): else: message = _FlatArrayMessageReader(buf, traversal_limit_in_words, nesting_limit) return message.get_root(self.schema) + def from_segments(self, segments): + message = _SegmentArrayMessageReader(segments) + return message.get_root(self.schema) def from_bytes_packed(self, buf, traversal_limit_in_words = None, nesting_limit = None): """Returns a Reader for the packed object in buf. @@ -3285,6 +3294,18 @@ cdef class _MessageBuilder: self.thisptr.setRootDynamicStruct((<_DynamicStructReader>value).thisptr) return self.get_root(value.schema) + cpdef get_segments_for_output(self) except +reraise_kj_exception: + segments = self.thisptr.getSegmentsForOutput() + res = [] + cdef const char* ptr + cdef bytes segment_bytes + for i in range(0, segments.size()): + segment = segments[i] + ptr = segment.begin() + segment_bytes = ptr[:8*segment.size()] + res.append(segment_bytes) + return res + cpdef new_orphan(self, schema) except +reraise_kj_exception: """A method for instantiating Cap'n Proto orphans @@ -3738,6 +3759,36 @@ cdef class _FlatArrayMessageReader(_MessageReader): del self.thisptr +@cython.internal +cdef class _SegmentArrayMessageReader(_MessageReader): + + cdef object _objects_to_pin + + def __init__(self, segments): + # takes a Python array of bytes and constructs a ConstWordArrayArrayPtr + num_segments = len(segments) + cdef char* ptr + cdef schema_cpp.ConstWordArrayPtr seg_ptr + cdef schema_cpp.ConstWordArrayPtr* seg_ptrs = malloc(num_segments * sizeof(schema_cpp.ConstWordArrayPtr)) + self._objects_to_pin = [] + for i in range(0, num_segments): + segment = bytes(segments[i]) + ptr = segment + if (ptr) % 8 != 0: + aligned = _AlignedBuffer(segment) + ptr = aligned.buf + self._objects_to_pin.append(aligned) + else: + self._objects_to_pin.append(segment) + seg_ptr = schema_cpp.ConstWordArrayPtr(ptr, len(segment)//8) + seg_ptrs[i] = seg_ptr + self.thisptr = new schema_cpp.SegmentArrayMessageReader( + schema_cpp.ConstWordArrayArrayPtr(seg_ptrs, num_segments)) + + def __dealloc__(self): + del self.thisptr + + @cython.internal cdef class _FlatMessageBuilder(_MessageBuilder): cdef object _object_to_pin From 73d6c5dc107c87d5b2b26460fc4486a5789f82b8 Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Thu, 23 Jun 2016 14:15:39 -0700 Subject: [PATCH 2/9] Add segment (de)serialization round-trip test. --- test/test_serialization.py | 7 +++++++ 1 file changed, 7 insertions(+) diff --git a/test/test_serialization.py b/test/test_serialization.py index d1b840c..5694088 100644 --- a/test/test_serialization.py +++ b/test/test_serialization.py @@ -42,6 +42,13 @@ def test_roundtrip_bytes(all_types): msg = all_types.TestAllTypes.from_bytes(message_bytes) test_regression.check_all_types(msg) +def test_roundtrip_segments(all_types): + msg = all_types.TestAllTypes.new_message() + test_regression.init_all_types(msg) + segments = msg.to_segments() + msg = all_types.TestAllTypes.from_segments(segments) + test_regression.check_all_types(msg) + @pytest.mark.skipif(sys.version_info[0] < 3, reason="mmap doesn't implement the buffer interface under python 2.") def test_roundtrip_bytes_mmap(all_types): msg = all_types.TestAllTypes.new_message() From 42ed1e819fbd28f296e4f8b8a15e1c781c7ed8fb Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Thu, 23 Jun 2016 14:15:46 -0700 Subject: [PATCH 3/9] Update documentation to cover segments. --- docs/quickstart.rst | 21 +++++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/docs/quickstart.rst b/docs/quickstart.rst index aa89b89..b19ad55 100644 --- a/docs/quickstart.rst +++ b/docs/quickstart.rst @@ -304,6 +304,27 @@ There are also packed versions:: alice2 = addressbook_capnp.Person.from_bytes_packed(alice.to_bytes_packed()) + +Byte Segments +~~~~~~~~~~~~~ + +Cap'n Proto supports a serialization mode which minimizes object copies. In the C++ interface, ``capnp::MessageBuilder::getSegmentsForOutput()`` returns an array of pointers to segments of the message's content without copying. ``capnp::SegmentArrayMessageReader`` performs the reverse operation, i.e., takes an array of pointers to segments and uses the underlying data, again without copying. This produces a different wire serialization format from ``to_bytes()`` serialization, which uses ``capnp::messageToFlatArray()`` and ``capnp::FlatArrayMessageReader`` (both of which use segments internally, but write them in an incompatible way). + +For compatibility on the Python side, use the ``to_segments()`` and ``from_segments()`` functions:: + + segments = alice.to_segments() + +This returns a list of segments, each a byte buffer. Each segment can be, e.g., turned into a ZeroMQ message frame. The list of segments can also be turned back into an object:: + + alice = addressbook_capnp.Person.from_segments(segments) + +For more information, please refer to the following links: + +- `Advice on minimizing copies from Cap'n Proto `_ (from the author of Cap'n Proto) +- `Advice on using Cap'n Proto over ZeroMQ `_ (from the author of Cap'n Proto) +- `Discussion about sending and reassembling Cap'n Proto message segments in C++ `_ (from the Cap'n Proto mailing list; includes sample code) + + RPC ---------- From 70e4e2d930a594733b148b901a8576635dc0c330 Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 13:25:27 -0700 Subject: [PATCH 4/9] Add ReaderOptions support to segment (de)serialization. --- capnp/lib/capnp.pyx | 16 +++++++++++----- 1 file changed, 11 insertions(+), 5 deletions(-) diff --git a/capnp/lib/capnp.pyx b/capnp/lib/capnp.pyx index 8ee1c3d..bc8c99b 100644 --- a/capnp/lib/capnp.pyx +++ b/capnp/lib/capnp.pyx @@ -2995,8 +2995,8 @@ class _StructModule(object): else: message = _FlatArrayMessageReader(buf, traversal_limit_in_words, nesting_limit) return message.get_root(self.schema) - def from_segments(self, segments): - message = _SegmentArrayMessageReader(segments) + def from_segments(self, segments, traversal_limit_in_words = None, nesting_limit = None): + message = _SegmentArrayMessageReader(segments, traversal_limit_in_words, nesting_limit) return message.get_root(self.schema) def from_bytes_packed(self, buf, traversal_limit_in_words = None, nesting_limit = None): """Returns a Reader for the packed object in buf. @@ -3764,8 +3764,13 @@ cdef class _SegmentArrayMessageReader(_MessageReader): cdef object _objects_to_pin - def __init__(self, segments): - # takes a Python array of bytes and constructs a ConstWordArrayArrayPtr + def __init__(self, segments, traversal_limit_in_words = None, nesting_limit = None): + cdef schema_cpp.ReaderOptions opts + if traversal_limit_in_words is not None: + opts.traversalLimitInWords = traversal_limit_in_words + if nesting_limit is not None: + opts.nestingLimit = nesting_limit + # take a Python array of bytes and constructs a ConstWordArrayArrayPtr num_segments = len(segments) cdef char* ptr cdef schema_cpp.ConstWordArrayPtr seg_ptr @@ -3783,7 +3788,8 @@ cdef class _SegmentArrayMessageReader(_MessageReader): seg_ptr = schema_cpp.ConstWordArrayPtr(ptr, len(segment)//8) seg_ptrs[i] = seg_ptr self.thisptr = new schema_cpp.SegmentArrayMessageReader( - schema_cpp.ConstWordArrayArrayPtr(seg_ptrs, num_segments)) + schema_cpp.ConstWordArrayArrayPtr(seg_ptrs, num_segments), + opts) def __dealloc__(self): del self.thisptr From fb8a3d8ac2622ee13e10c502eb00b23219dfb4aa Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 13:27:21 -0700 Subject: [PATCH 5/9] Add docstrings to (de)serialization functions. --- capnp/lib/capnp.pyx | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/capnp/lib/capnp.pyx b/capnp/lib/capnp.pyx index bc8c99b..263515d 100644 --- a/capnp/lib/capnp.pyx +++ b/capnp/lib/capnp.pyx @@ -1193,6 +1193,14 @@ cdef class _DynamicStructBuilder: return ret cpdef to_segments(_DynamicStructBuilder self) except +reraise_kj_exception: + """Returns the struct's containing message as a Python list of Python bytes objects. + + This avoids making copies. + + NB: This is not currently supported on PyPy. + + :rtype: list + """ self._check_write() cdef _MessageBuilder builder = self._parent segments = builder.get_segments_for_output() @@ -2996,6 +3004,14 @@ class _StructModule(object): message = _FlatArrayMessageReader(buf, traversal_limit_in_words, nesting_limit) return message.get_root(self.schema) def from_segments(self, segments, traversal_limit_in_words = None, nesting_limit = None): + """Returns a Reader for a list of segment bytes. + + This avoids making copies. + + NB: This is not currently supported on PyPy. + + :rtype: list + """ message = _SegmentArrayMessageReader(segments, traversal_limit_in_words, nesting_limit) return message.get_root(self.schema) def from_bytes_packed(self, buf, traversal_limit_in_words = None, nesting_limit = None): From 5a63dacf146cafe8614596b79ae7b870bcd7fcef Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 13:27:39 -0700 Subject: [PATCH 6/9] Disable segment serialization test on PyPy (for now). --- test/test_serialization.py | 1 + 1 file changed, 1 insertion(+) diff --git a/test/test_serialization.py b/test/test_serialization.py index 5694088..ebead7e 100644 --- a/test/test_serialization.py +++ b/test/test_serialization.py @@ -42,6 +42,7 @@ def test_roundtrip_bytes(all_types): msg = all_types.TestAllTypes.from_bytes(message_bytes) test_regression.check_all_types(msg) +@pytest.mark.skipif(platform.python_implementation() == 'PyPy', reason="TODO: Investigate why this works on CPython but fails on PyPy.") def test_roundtrip_segments(all_types): msg = all_types.TestAllTypes.new_message() test_regression.init_all_types(msg) From f703a7f099ed24c783a24fc889a41c8af27734f0 Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 13:27:54 -0700 Subject: [PATCH 7/9] Update segment (de)serialization documentation. --- docs/quickstart.rst | 2 ++ 1 file changed, 2 insertions(+) diff --git a/docs/quickstart.rst b/docs/quickstart.rst index b19ad55..3b7a1a8 100644 --- a/docs/quickstart.rst +++ b/docs/quickstart.rst @@ -308,6 +308,8 @@ There are also packed versions:: Byte Segments ~~~~~~~~~~~~~ +.. note:: This feature is not supported in PyPy at the moment, pending investigation. + Cap'n Proto supports a serialization mode which minimizes object copies. In the C++ interface, ``capnp::MessageBuilder::getSegmentsForOutput()`` returns an array of pointers to segments of the message's content without copying. ``capnp::SegmentArrayMessageReader`` performs the reverse operation, i.e., takes an array of pointers to segments and uses the underlying data, again without copying. This produces a different wire serialization format from ``to_bytes()`` serialization, which uses ``capnp::messageToFlatArray()`` and ``capnp::FlatArrayMessageReader`` (both of which use segments internally, but write them in an incompatible way). For compatibility on the Python side, use the ``to_segments()`` and ``from_segments()`` functions:: From e8a8d26260cba0fb69670bc474e8db68990259b7 Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 14:03:20 -0700 Subject: [PATCH 8/9] Fix memory leak. --- capnp/lib/capnp.pyx | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/capnp/lib/capnp.pyx b/capnp/lib/capnp.pyx index 263515d..54344b2 100644 --- a/capnp/lib/capnp.pyx +++ b/capnp/lib/capnp.pyx @@ -3779,6 +3779,7 @@ cdef class _FlatArrayMessageReader(_MessageReader): cdef class _SegmentArrayMessageReader(_MessageReader): cdef object _objects_to_pin + cdef schema_cpp.ConstWordArrayPtr* _seg_ptrs def __init__(self, segments, traversal_limit_in_words = None, nesting_limit = None): cdef schema_cpp.ReaderOptions opts @@ -3790,7 +3791,7 @@ cdef class _SegmentArrayMessageReader(_MessageReader): num_segments = len(segments) cdef char* ptr cdef schema_cpp.ConstWordArrayPtr seg_ptr - cdef schema_cpp.ConstWordArrayPtr* seg_ptrs = malloc(num_segments * sizeof(schema_cpp.ConstWordArrayPtr)) + self._seg_ptrs = malloc(num_segments * sizeof(schema_cpp.ConstWordArrayPtr)) self._objects_to_pin = [] for i in range(0, num_segments): segment = bytes(segments[i]) @@ -3802,12 +3803,13 @@ cdef class _SegmentArrayMessageReader(_MessageReader): else: self._objects_to_pin.append(segment) seg_ptr = schema_cpp.ConstWordArrayPtr(ptr, len(segment)//8) - seg_ptrs[i] = seg_ptr + self._seg_ptrs[i] = seg_ptr self.thisptr = new schema_cpp.SegmentArrayMessageReader( - schema_cpp.ConstWordArrayArrayPtr(seg_ptrs, num_segments), + schema_cpp.ConstWordArrayArrayPtr(self._seg_ptrs, num_segments), opts) def __dealloc__(self): + free(self._seg_ptrs) del self.thisptr From 41c418aa1cbc2a7d130176ab08326f3f488238ac Mon Sep 17 00:00:00 2001 From: Constantine Vetoshev Date: Fri, 24 Jun 2016 16:43:57 -0700 Subject: [PATCH 9/9] Use PyObject_AsReadBuffer instead of bytes(). --- capnp/lib/capnp.pyx | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/capnp/lib/capnp.pyx b/capnp/lib/capnp.pyx index 54344b2..5d44e8b 100644 --- a/capnp/lib/capnp.pyx +++ b/capnp/lib/capnp.pyx @@ -3789,20 +3789,20 @@ cdef class _SegmentArrayMessageReader(_MessageReader): opts.nestingLimit = nesting_limit # take a Python array of bytes and constructs a ConstWordArrayArrayPtr num_segments = len(segments) - cdef char* ptr + cdef const void* ptr + cdef Py_ssize_t segment_size cdef schema_cpp.ConstWordArrayPtr seg_ptr self._seg_ptrs = malloc(num_segments * sizeof(schema_cpp.ConstWordArrayPtr)) self._objects_to_pin = [] for i in range(0, num_segments): - segment = bytes(segments[i]) - ptr = segment + PyObject_AsReadBuffer(segments[i], &ptr, &segment_size) if (ptr) % 8 != 0: - aligned = _AlignedBuffer(segment) + aligned = _AlignedBuffer(segments[i]) ptr = aligned.buf self._objects_to_pin.append(aligned) else: - self._objects_to_pin.append(segment) - seg_ptr = schema_cpp.ConstWordArrayPtr(ptr, len(segment)//8) + self._objects_to_pin.append(segments[i]) + seg_ptr = schema_cpp.ConstWordArrayPtr(ptr, segment_size//8) self._seg_ptrs[i] = seg_ptr self.thisptr = new schema_cpp.SegmentArrayMessageReader( schema_cpp.ConstWordArrayArrayPtr(self._seg_ptrs, num_segments),