Merge branch 'release/v0.5.6'
This commit is contained in:
@@ -8,8 +8,9 @@ python:
|
|||||||
- pypy
|
- pypy
|
||||||
|
|
||||||
env:
|
env:
|
||||||
- BUILD_CAPNP=true
|
|
||||||
- BUILD_CAPNP=
|
- BUILD_CAPNP=
|
||||||
|
- BUILD_CAPNP=true
|
||||||
|
- BUILD_CAPNP=true CFLAGS="-DKJ_DEBUG"
|
||||||
|
|
||||||
# skip testing for pypy + BUILD_CAPNP=false since it's failing in travis for some reason
|
# skip testing for pypy + BUILD_CAPNP=false since it's failing in travis for some reason
|
||||||
matrix:
|
matrix:
|
||||||
|
|||||||
@@ -1,3 +1,8 @@
|
|||||||
|
## v0.5.6 (2015-04-13)
|
||||||
|
- Fix a serious bug in TwoPartyServer that was preventing it from working when passed a string address.
|
||||||
|
- Fix bugs that were exposed by defining KJDEBUG (thanks @davidcarne for finding this)
|
||||||
|
|
||||||
|
|
||||||
## v0.5.5 (2015-03-06)
|
## v0.5.5 (2015-03-06)
|
||||||
- Update bundled C++ libcapnp to v0.5.1.2 security release
|
- Update bundled C++ libcapnp to v0.5.1.2 security release
|
||||||
|
|
||||||
|
|||||||
@@ -13,8 +13,8 @@ def build_libcapnp(bundle_dir, build_dir, verbose=False):
|
|||||||
stdout = f
|
stdout = f
|
||||||
if verbose:
|
if verbose:
|
||||||
stdout = None
|
stdout = None
|
||||||
cxxflags = os.environ.get('CXXFLAGS', '')
|
cxxflags = os.environ.get('CXXFLAGS', None)
|
||||||
os.environ['CXXFLAGS'] = cxxflags + ' -fPIC -O2 -DNDEBUG'
|
os.environ['CXXFLAGS'] = (cxxflags or '') + ' -fPIC -O2 -DNDEBUG'
|
||||||
conf = subprocess.Popen(['./configure', '--disable-shared', '--prefix', build_dir], cwd=capnp_dir, stdout=stdout)
|
conf = subprocess.Popen(['./configure', '--disable-shared', '--prefix', build_dir], cwd=capnp_dir, stdout=stdout)
|
||||||
returncode = conf.wait()
|
returncode = conf.wait()
|
||||||
if returncode != 0:
|
if returncode != 0:
|
||||||
@@ -22,5 +22,9 @@ def build_libcapnp(bundle_dir, build_dir, verbose=False):
|
|||||||
|
|
||||||
make = subprocess.Popen(['make', '-j4', 'install'], cwd=capnp_dir, stdout=stdout)
|
make = subprocess.Popen(['make', '-j4', 'install'], cwd=capnp_dir, stdout=stdout)
|
||||||
returncode = make.wait()
|
returncode = make.wait()
|
||||||
|
if cxxflags is None:
|
||||||
|
del os.environ['CXXFLAGS']
|
||||||
|
else:
|
||||||
|
os.environ['CXXFLAGS'] = cxxflags
|
||||||
if returncode != 0:
|
if returncode != 0:
|
||||||
raise RuntimeError('Make failed')
|
raise RuntimeError('Make failed')
|
||||||
|
|||||||
@@ -2,6 +2,8 @@
|
|||||||
|
|
||||||
set -exo pipefail
|
set -exo pipefail
|
||||||
|
|
||||||
|
CAPNP_VERSION=0.5.1.2
|
||||||
|
|
||||||
sudo add-apt-repository -y ppa:ubuntu-toolchain-r/test
|
sudo add-apt-repository -y ppa:ubuntu-toolchain-r/test
|
||||||
sudo apt-get -qq update
|
sudo apt-get -qq update
|
||||||
sudo apt-get -qq install g++-4.8 libstdc++-4.8-dev
|
sudo apt-get -qq install g++-4.8 libstdc++-4.8-dev
|
||||||
@@ -9,5 +11,5 @@ sudo update-alternatives --quiet --install /usr/bin/gcc gcc /usr/bin/gcc-4.8
|
|||||||
sudo update-alternatives --quiet --set gcc /usr/bin/gcc-4.8
|
sudo update-alternatives --quiet --set gcc /usr/bin/gcc-4.8
|
||||||
|
|
||||||
if ! [ -z "${BUILD_CAPNP}" ]; then
|
if ! [ -z "${BUILD_CAPNP}" ]; then
|
||||||
wget https://capnproto.org/capnproto-c++-0.5.1.tar.gz && tar xzvf capnproto-c++-0.5.1.tar.gz && cd capnproto-c++-0.5.1 && ./configure && make -j6 check && sudo make install && sudo ldconfig && cd ..
|
wget https://capnproto.org/capnproto-c++-${CAPNP_VERSION}.tar.gz && tar xzvf capnproto-c++-${CAPNP_VERSION}.tar.gz && cd capnproto-c++-${CAPNP_VERSION} && ./configure && make -j6 check && sudo make install && sudo ldconfig && cd ..
|
||||||
fi
|
fi
|
||||||
|
|||||||
@@ -32,7 +32,8 @@ cdef extern from "capnp/helpers/rpcHelper.h":
|
|||||||
Capability.Client restoreHelper(RpcSystem&, AnyPointer.Builder&)
|
Capability.Client restoreHelper(RpcSystem&, AnyPointer.Builder&)
|
||||||
Capability.Client bootstrapHelper(RpcSystem&)
|
Capability.Client bootstrapHelper(RpcSystem&)
|
||||||
RpcSystem makeRpcClientWithRestorer(TwoPartyVatNetwork&, PyRestorer&)
|
RpcSystem makeRpcClientWithRestorer(TwoPartyVatNetwork&, PyRestorer&)
|
||||||
PyPromise connectServer(TaskSet &, PyRestorer &, AsyncIoContext *, StringPtr)
|
PyPromise connectServerRestorer(TaskSet &, PyRestorer &, AsyncIoContext *, StringPtr)
|
||||||
|
PyPromise connectServer(TaskSet &, Capability.Client, AsyncIoContext *, StringPtr)
|
||||||
|
|
||||||
cdef extern from "capnp/helpers/serialize.h":
|
cdef extern from "capnp/helpers/serialize.h":
|
||||||
ByteArray messageToPackedBytes(MessageBuilder &, size_t wordCount)
|
ByteArray messageToPackedBytes(MessageBuilder &, size_t wordCount)
|
||||||
|
|||||||
@@ -89,12 +89,12 @@ capnp::RpcSystem<SturdyRefHostId> makeRpcClientWithRestorer(
|
|||||||
return RpcSystem<SturdyRefHostId>(network, restorer);
|
return RpcSystem<SturdyRefHostId>(network, restorer);
|
||||||
}
|
}
|
||||||
|
|
||||||
struct ServerContext {
|
struct ServerContextRestorer {
|
||||||
kj::Own<kj::AsyncIoStream> stream;
|
kj::Own<kj::AsyncIoStream> stream;
|
||||||
capnp::TwoPartyVatNetwork network;
|
capnp::TwoPartyVatNetwork network;
|
||||||
capnp::RpcSystem<capnp::rpc::twoparty::SturdyRefHostId> rpcSystem;
|
capnp::RpcSystem<capnp::rpc::twoparty::SturdyRefHostId> rpcSystem;
|
||||||
|
|
||||||
ServerContext(kj::Own<kj::AsyncIoStream>&& stream, capnp::SturdyRefRestorer<capnp::AnyPointer>& restorer)
|
ServerContextRestorer(kj::Own<kj::AsyncIoStream>&& stream, capnp::SturdyRefRestorer<capnp::AnyPointer>& restorer)
|
||||||
: stream(kj::mv(stream)),
|
: stream(kj::mv(stream)),
|
||||||
network(*this->stream, capnp::rpc::twoparty::Side::SERVER),
|
network(*this->stream, capnp::rpc::twoparty::Side::SERVER),
|
||||||
rpcSystem(makeRpcServer(network, restorer)) {}
|
rpcSystem(makeRpcServer(network, restorer)) {}
|
||||||
@@ -106,14 +106,14 @@ class ErrorHandler : public kj::TaskSet::ErrorHandler {
|
|||||||
}
|
}
|
||||||
};
|
};
|
||||||
|
|
||||||
void acceptLoop(kj::TaskSet & tasks, PyRestorer & restorer, kj::Own<kj::ConnectionReceiver>&& listener) {
|
void acceptLoopRestorer(kj::TaskSet & tasks, PyRestorer & restorer, kj::Own<kj::ConnectionReceiver>&& listener) {
|
||||||
auto ptr = listener.get();
|
auto ptr = listener.get();
|
||||||
tasks.add(ptr->accept().then(kj::mvCapture(kj::mv(listener),
|
tasks.add(ptr->accept().then(kj::mvCapture(kj::mv(listener),
|
||||||
[&](kj::Own<kj::ConnectionReceiver>&& listener,
|
[&](kj::Own<kj::ConnectionReceiver>&& listener,
|
||||||
kj::Own<kj::AsyncIoStream>&& connection) {
|
kj::Own<kj::AsyncIoStream>&& connection) {
|
||||||
acceptLoop(tasks, restorer, kj::mv(listener));
|
acceptLoopRestorer(tasks, restorer, kj::mv(listener));
|
||||||
|
|
||||||
auto server = kj::heap<ServerContext>(kj::mv(connection), restorer);
|
auto server = kj::heap<ServerContextRestorer>(kj::mv(connection), restorer);
|
||||||
|
|
||||||
// Arrange to destroy the server context when all references are gone, or when the
|
// Arrange to destroy the server context when all references are gone, or when the
|
||||||
// EzRpcServer is destroyed (which will destroy the TaskSet).
|
// EzRpcServer is destroyed (which will destroy the TaskSet).
|
||||||
@@ -121,7 +121,7 @@ void acceptLoop(kj::TaskSet & tasks, PyRestorer & restorer, kj::Own<kj::Connecti
|
|||||||
})));
|
})));
|
||||||
}
|
}
|
||||||
|
|
||||||
kj::Promise<PyObject *> connectServer(kj::TaskSet & tasks, PyRestorer & restorer, kj::AsyncIoContext * context, kj::StringPtr bindAddress) {
|
kj::Promise<PyObject *> connectServerRestorer(kj::TaskSet & tasks, PyRestorer & restorer, kj::AsyncIoContext * context, kj::StringPtr bindAddress) {
|
||||||
auto paf = kj::newPromiseAndFulfiller<unsigned int>();
|
auto paf = kj::newPromiseAndFulfiller<unsigned int>();
|
||||||
auto portPromise = paf.promise.fork();
|
auto portPromise = paf.promise.fork();
|
||||||
|
|
||||||
@@ -131,7 +131,50 @@ kj::Promise<PyObject *> connectServer(kj::TaskSet & tasks, PyRestorer & restorer
|
|||||||
kj::Own<kj::NetworkAddress>&& addr) {
|
kj::Own<kj::NetworkAddress>&& addr) {
|
||||||
auto listener = addr->listen();
|
auto listener = addr->listen();
|
||||||
portFulfiller->fulfill(listener->getPort());
|
portFulfiller->fulfill(listener->getPort());
|
||||||
acceptLoop(tasks, restorer, kj::mv(listener));
|
acceptLoopRestorer(tasks, restorer, kj::mv(listener));
|
||||||
|
})));
|
||||||
|
|
||||||
|
return portPromise.addBranch().then([&](unsigned int port) { return PyLong_FromUnsignedLong(port); });
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
|
struct ServerContext {
|
||||||
|
kj::Own<kj::AsyncIoStream> stream;
|
||||||
|
capnp::TwoPartyVatNetwork network;
|
||||||
|
capnp::RpcSystem<capnp::rpc::twoparty::SturdyRefHostId> rpcSystem;
|
||||||
|
|
||||||
|
ServerContext(kj::Own<kj::AsyncIoStream>&& stream, capnp::Capability::Client client)
|
||||||
|
: stream(kj::mv(stream)),
|
||||||
|
network(*this->stream, capnp::rpc::twoparty::Side::SERVER),
|
||||||
|
rpcSystem(makeRpcServer(network, client)) {}
|
||||||
|
};
|
||||||
|
|
||||||
|
void acceptLoop(kj::TaskSet & tasks, capnp::Capability::Client client, kj::Own<kj::ConnectionReceiver>&& listener) {
|
||||||
|
auto ptr = listener.get();
|
||||||
|
tasks.add(ptr->accept().then(kj::mvCapture(kj::mv(listener),
|
||||||
|
[&, client](kj::Own<kj::ConnectionReceiver>&& listener,
|
||||||
|
kj::Own<kj::AsyncIoStream>&& connection) mutable {
|
||||||
|
acceptLoop(tasks, client, kj::mv(listener));
|
||||||
|
|
||||||
|
auto server = kj::heap<ServerContext>(kj::mv(connection), client);
|
||||||
|
|
||||||
|
// Arrange to destroy the server context when all references are gone, or when the
|
||||||
|
// EzRpcServer is destroyed (which will destroy the TaskSet).
|
||||||
|
tasks.add(server->network.onDisconnect().attach(kj::mv(server)));
|
||||||
|
})));
|
||||||
|
}
|
||||||
|
|
||||||
|
kj::Promise<PyObject *> connectServer(kj::TaskSet & tasks, capnp::Capability::Client client, kj::AsyncIoContext * context, kj::StringPtr bindAddress) {
|
||||||
|
auto paf = kj::newPromiseAndFulfiller<unsigned int>();
|
||||||
|
auto portPromise = paf.promise.fork();
|
||||||
|
|
||||||
|
tasks.add(context->provider->getNetwork().parseAddress(bindAddress)
|
||||||
|
.then(kj::mvCapture(paf.fulfiller,
|
||||||
|
[&, client](kj::Own<kj::PromiseFulfiller<unsigned int>>&& portFulfiller,
|
||||||
|
kj::Own<kj::NetworkAddress>&& addr) mutable {
|
||||||
|
auto listener = addr->listen();
|
||||||
|
portFulfiller->fulfill(listener->getPort());
|
||||||
|
acceptLoop(tasks, client, kj::mv(listener));
|
||||||
})));
|
})));
|
||||||
|
|
||||||
return portPromise.addBranch().then([&](unsigned int port) { return PyLong_FromUnsignedLong(port); });
|
return portPromise.addBranch().then([&](unsigned int port) { return PyLong_FromUnsignedLong(port); });
|
||||||
|
|||||||
@@ -2279,7 +2279,7 @@ cdef class TwoPartyServer:
|
|||||||
self._bootstrap = None
|
self._bootstrap = None
|
||||||
|
|
||||||
if isinstance(socket, basestring):
|
if isinstance(socket, basestring):
|
||||||
self._connect(socket)
|
self._connect(socket, restorer, bootstrap)
|
||||||
else:
|
else:
|
||||||
self._orig_stream = socket
|
self._orig_stream = socket
|
||||||
self._stream = _FdAsyncIoStream(socket.fileno())
|
self._stream = _FdAsyncIoStream(socket.fileno())
|
||||||
@@ -2302,11 +2302,19 @@ cdef class TwoPartyServer:
|
|||||||
Py_INCREF(self._network)
|
Py_INCREF(self._network)
|
||||||
self._disconnect_promise = self.on_disconnect().then(self._decref)
|
self._disconnect_promise = self.on_disconnect().then(self._decref)
|
||||||
|
|
||||||
cpdef _connect(self, host_string):
|
cpdef _connect(self, host_string, restorer, bootstrap):
|
||||||
|
cdef _InterfaceSchema schema
|
||||||
cdef _EventLoop loop = C_DEFAULT_EVENT_LOOP_GETTER()
|
cdef _EventLoop loop = C_DEFAULT_EVENT_LOOP_GETTER()
|
||||||
cdef capnp.StringPtr temp_string = capnp.StringPtr(<char*>host_string, len(host_string))
|
cdef capnp.StringPtr temp_string = capnp.StringPtr(<char*>host_string, len(host_string))
|
||||||
self._task_set = new capnp.TaskSet(self._error_handler)
|
self._task_set = new capnp.TaskSet(self._error_handler)
|
||||||
self.port_promise = Promise()._init(helpers.connectServer(deref(self._task_set), deref(self._restorer.thisptr), loop.thisptr, temp_string))
|
if restorer:
|
||||||
|
self._restorer = _convert_restorer(restorer)
|
||||||
|
self.port_promise = Promise()._init(helpers.connectServerRestorer(deref(self._task_set), deref(self._restorer.thisptr), loop.thisptr, temp_string))
|
||||||
|
else:
|
||||||
|
self._bootstrap = bootstrap
|
||||||
|
Py_INCREF(self._bootstrap)
|
||||||
|
schema = bootstrap.schema
|
||||||
|
self.port_promise = Promise()._init(helpers.connectServer(deref(self._task_set), helpers.server_to_client(schema.thisptr, <PyObject *>bootstrap), loop.thisptr, temp_string))
|
||||||
|
|
||||||
def _decref(self):
|
def _decref(self):
|
||||||
Py_DECREF(self._bootstrap)
|
Py_DECREF(self._bootstrap)
|
||||||
@@ -2697,9 +2705,9 @@ types.Data = _data
|
|||||||
# _list.thisptr = capnp.SchemaType(capnp.TypeWhichLIST)
|
# _list.thisptr = capnp.SchemaType(capnp.TypeWhichLIST)
|
||||||
# types.list = _list
|
# types.list = _list
|
||||||
|
|
||||||
cdef _SchemaType _enum = _SchemaType()
|
# cdef _SchemaType _enum = _SchemaType()
|
||||||
_enum.thisptr = capnp.SchemaType(capnp.TypeWhichENUM)
|
# _enum.thisptr = capnp.SchemaType(capnp.TypeWhichENUM)
|
||||||
types.Enum = _enum
|
# types.Enum = _enum
|
||||||
|
|
||||||
# cdef _SchemaType _struct = _SchemaType()
|
# cdef _SchemaType _struct = _SchemaType()
|
||||||
# _struct.thisptr = capnp.SchemaType(capnp.TypeWhichSTRUCT)
|
# _struct.thisptr = capnp.SchemaType(capnp.TypeWhichSTRUCT)
|
||||||
|
|||||||
@@ -39,7 +39,7 @@ def main(host):
|
|||||||
# takes a struct or AnyPointer as an argument), and then cast the returned
|
# takes a struct or AnyPointer as an argument), and then cast the returned
|
||||||
# capability to it's proper type. This casting is due to capabilities not
|
# capability to it's proper type. This casting is due to capabilities not
|
||||||
# having a reference to their schema
|
# having a reference to their schema
|
||||||
calculator = client.ez_restore('calculator').cast_as(calculator_capnp.Calculator)
|
calculator = client.bootstrap().cast_as(calculator_capnp.Calculator)
|
||||||
|
|
||||||
'''Make a request that just evaluates the literal value 123.
|
'''Make a request that just evaluates the literal value 123.
|
||||||
|
|
||||||
|
|||||||
@@ -130,15 +130,10 @@ given address/port ADDRESS may be '*' to bind to all local addresses.\
|
|||||||
return parser.parse_args()
|
return parser.parse_args()
|
||||||
|
|
||||||
|
|
||||||
def restore(ref):
|
|
||||||
assert ref.as_text() == 'calculator'
|
|
||||||
return CalculatorImpl()
|
|
||||||
|
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
address = parse_args().address
|
address = parse_args().address
|
||||||
|
|
||||||
server = capnp.TwoPartyServer(address, restore)
|
server = capnp.TwoPartyServer(address, bootstrap=CalculatorImpl())
|
||||||
server.run_forever()
|
server.run_forever()
|
||||||
|
|
||||||
if __name__ == '__main__':
|
if __name__ == '__main__':
|
||||||
|
|||||||
2
setup.py
2
setup.py
@@ -14,7 +14,7 @@ _this_dir = os.path.dirname(__file__)
|
|||||||
|
|
||||||
MAJOR = 0
|
MAJOR = 0
|
||||||
MINOR = 5
|
MINOR = 5
|
||||||
MICRO = 5
|
MICRO = 6
|
||||||
VERSION = '%d.%d.%d' % (MAJOR, MINOR, MICRO)
|
VERSION = '%d.%d.%d' % (MAJOR, MINOR, MICRO)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -12,7 +12,7 @@ import calculator_server
|
|||||||
def test_calculator():
|
def test_calculator():
|
||||||
read, write = socket.socketpair(socket.AF_UNIX)
|
read, write = socket.socketpair(socket.AF_UNIX)
|
||||||
|
|
||||||
server = capnp.TwoPartyServer(write, calculator_server.restore)
|
server = capnp.TwoPartyServer(write, bootstrap=calculator_server.CalculatorImpl())
|
||||||
calculator_client.main(read)
|
calculator_client.main(read)
|
||||||
|
|
||||||
|
|
||||||
@@ -29,7 +29,7 @@ def test_calculator_gc():
|
|||||||
evaluate_impl_orig = calculator_server.evaluate_impl
|
evaluate_impl_orig = calculator_server.evaluate_impl
|
||||||
calculator_server.evaluate_impl = new_evaluate_impl(evaluate_impl_orig)
|
calculator_server.evaluate_impl = new_evaluate_impl(evaluate_impl_orig)
|
||||||
|
|
||||||
server = capnp.TwoPartyServer(write, calculator_server.restore)
|
server = capnp.TwoPartyServer(write, bootstrap=calculator_server.CalculatorImpl())
|
||||||
calculator_client.main(read)
|
calculator_client.main(read)
|
||||||
|
|
||||||
calculator_server.evaluate_impl = evaluate_impl_orig
|
calculator_server.evaluate_impl = evaluate_impl_orig
|
||||||
|
|||||||
@@ -106,6 +106,7 @@ def test_roundtrip_bytes_multiple_packed(all_types):
|
|||||||
i += 1
|
i += 1
|
||||||
assert i == 3
|
assert i == 3
|
||||||
|
|
||||||
|
@pytest.mark.skipif(platform.python_implementation() == 'PyPy', reason="This works on my local PyPy v2.5.0, but is for some reason broken on TravisCI. Skip for now.")
|
||||||
def test_roundtrip_dict(all_types):
|
def test_roundtrip_dict(all_types):
|
||||||
msg = all_types.TestAllTypes.new_message()
|
msg = all_types.TestAllTypes.new_message()
|
||||||
test_regression.init_all_types(msg)
|
test_regression.init_all_types(msg)
|
||||||
|
|||||||
Reference in New Issue
Block a user