5#include <mp/test/foo.capnp.h>
6#include <mp/test/foo.capnp.proxy.h>
9#include <capnp/capability.h>
13#include <condition_variable>
19#include <kj/async-io.h>
21#include <kj/exception.h>
29#include <mp/proxy.capnp.h>
40#include <unordered_set>
45#define EXPECT_EXCEPTION(call, message) \
49 } catch (const std::runtime_error& e) { \
50 KJ_EXPECT(std::string_view{e.what()} == message); \
59static_assert(std::is_integral_v<
decltype(
kMP_MAJOR_VERSION)>,
"MP_MAJOR_VERSION must be an integral constant");
60static_assert(std::is_integral_v<
decltype(
kMP_MINOR_VERSION)>,
"MP_MINOR_VERSION must be an integral constant");
82 std::promise<std::unique_ptr<ProxyClient<messages::FooInterface>>>
client_promise;
83 std::unique_ptr<ProxyClient<messages::FooInterface>>
client;
98 auto server_connection =
99 std::make_unique<Connection>(loop, kj::mv(pipe.ends[0]), [&](
Connection& connection) {
100 auto server_proxy = kj::heap<ProxyServer<messages::FooInterface>>(
101 std::make_shared<FooImplementation>(), connection);
102 server = server_proxy;
103 return capnp::Capability::Client(kj::mv(server_proxy));
107 assert(std::this_thread::get_id() == loop.m_thread_id);
108 loop.m_task_set->add(kj::evalLater([&] { server_connection.reset(); }));
112 server_connection->onDisconnect([&] { server_connection.reset(); });
114 auto client_connection = std::make_unique<Connection>(loop, kj::mv(pipe.ends[1]));
115 auto client_proxy = std::make_unique<ProxyClient<messages::FooInterface>>(
116 client_connection->m_rpc_system->bootstrap(
ServerVatId().vat_id).castAs<messages::FooInterface>(),
117 client_connection.get(), client_owns_connection);
118 if (client_owns_connection) {
119 (void)client_connection.release();
134 bool destroyed =
false;
135 client->m_context.cleanup_fns.emplace_front([&destroyed] { destroyed =
true; });
137 KJ_EXPECT(destroyed);
148 KJ_EXPECT(foo->add(1, 2) == 3);
150 foo->addOut(3, 4,
ret);
152 foo->addInOut(3,
ret);
153 KJ_EXPECT(
ret == 10);
168 KJ_EXPECT(in.
name ==
out.name);
169 KJ_EXPECT(in.
set_int.size() ==
out.set_int.size());
171 KJ_EXPECT(*
init == *outit);
175 KJ_EXPECT(
out.unordered_set_int.count(elem) == 1);
178 for (
size_t i = 0; i < in.
vector_bool.size(); ++i) {
184 KJ_EXPECT(
init->first == outit->first);
185 KJ_EXPECT(
init->second == outit->second);
189 KJ_EXPECT(foo->pass(in).optional_int == 3);
191 KJ_EXPECT(!foo->pass(in).optional_int);
205 int call(
int arg)
override
207 KJ_EXPECT(arg == m_expect);
210 int callExtended(
int arg)
override
212 KJ_EXPECT(arg == m_expect + 10);
218 foo->initThreadMap();
219 Callback callback(1, 2);
220 KJ_EXPECT(foo->callback(callback, 1) == 2);
221 KJ_EXPECT(foo->callbackUnique(std::make_unique<Callback>(3, 4), 3) == 4);
222 KJ_EXPECT(foo->callbackShared(std::make_shared<Callback>(5, 6), 5) == 6);
223 auto saved = std::make_shared<Callback>(7, 8);
224 KJ_EXPECT(saved.use_count() == 1);
225 foo->saveCallback(saved);
226 KJ_EXPECT(saved.use_count() == 2);
227 foo->callbackSaved(7);
228 KJ_EXPECT(foo->callbackSaved(7) == 8);
229 foo->saveCallback(
nullptr);
230 KJ_EXPECT(saved.use_count() == 1);
231 KJ_EXPECT(foo->callbackExtended(callback, 11) == 12);
236 FooCustom custom_out = foo->passCustom(custom_in);
237 KJ_EXPECT(custom_in.
v1 == custom_out.
v1);
238 KJ_EXPECT(custom_in.
v2 == custom_out.
v2);
243 KJ_EXPECT(empty_data_out.empty());
247 FooMessage message2{foo->passMessage(message1)};
248 KJ_EXPECT(message2.message ==
"init build read call build read");
252 foo->passMutable(mut);
253 KJ_EXPECT(mut.
message ==
"init build pass call return read");
255 KJ_EXPECT(foo->passDouble(1.25) == 1.25);
257 KJ_EXPECT(foo->passFn([]{ return 10; }) == 10);
260 KJ_EXPECT(foo->passFn([foo]{
261 return foo->passFn([]{ return 1; });
264 std::vector<FooDataRef> data_in;
265 data_in.push_back(std::make_shared<FooData>(
FooData{
'H',
'i'}));
266 data_in.push_back(
nullptr);
267 std::vector<FooDataRef> data_out{foo->passDataPointers(data_in)};
268 KJ_EXPECT(data_out.size() == 2);
269 KJ_REQUIRE(data_out[0] !=
nullptr);
270 KJ_EXPECT(*data_out[0] == *data_in[0]);
271 KJ_EXPECT(!data_out[1]);
274KJ_TEST(
"Call IPC method after client connection is closed")
278 KJ_EXPECT(foo->add(1, 2) == 3);
279 setup.client_disconnect();
281 EXPECT_EXCEPTION(foo->add(1, 2),
"IPC client method called after disconnect.");
284KJ_TEST(
"Calling IPC method after server connection is closed")
288 KJ_EXPECT(foo->add(1, 2) == 3);
289 setup.server_disconnect();
291 EXPECT_EXCEPTION(foo->add(1, 2),
"IPC client method call interrupted by disconnect.");
294KJ_TEST(
"Calling IPC method and disconnecting during the call")
298 KJ_EXPECT(foo->add(1, 2) == 3);
302 setup.server->m_impl->m_fn =
setup.client_disconnect;
304 EXPECT_EXCEPTION(foo->callFn(),
"IPC client method call interrupted by disconnect.");
307KJ_TEST(
"Calling IPC method, disconnecting and blocking during the call")
327 std::promise<void> signal;
330 KJ_EXPECT(foo->add(1, 2) == 3);
332 foo->initThreadMap();
333 setup.server->m_impl->m_fn = [&] {
335 setup.client_disconnect();
336 signal.get_future().get();
339 EXPECT_EXCEPTION(foo->callFnAsync(),
"IPC client method call interrupted by disconnect.");
349KJ_TEST(
"Worker thread destroyed before it is initialized")
360 foo->initThreadMap();
361 setup.server->m_impl->m_fn = [] {};
368 setup.server_disconnect_later();
373 std::this_thread::sleep_for(std::chrono::milliseconds(10));
376 EXPECT_EXCEPTION(foo->callFnAsync(),
"IPC client method call interrupted by disconnect.");
379KJ_TEST(
"Calling async IPC method, with server disconnect racing the call")
390 foo->initThreadMap();
391 setup.server->m_impl->m_fn = [] {};
395 setup.server_disconnect();
398 std::this_thread::sleep_for(std::chrono::milliseconds(10));
401 EXPECT_EXCEPTION(foo->callFnAsync(),
"IPC client method call interrupted by disconnect.");
404KJ_TEST(
"Calling async IPC method, with server disconnect after cleanup")
417 foo->initThreadMap();
418 setup.server->m_impl->m_fn = [] {};
422 setup.server_disconnect();
425 EXPECT_EXCEPTION(foo->callFnAsync(),
"IPC client method call interrupted by disconnect.");
428KJ_TEST(
"Destroying ProxyClient<> with destroy method after peer disconnect")
442 foo->initThreadMap();
447 int call(
int arg)
override {
return arg; }
450 foo->saveCallback(std::make_shared<Callback>());
451 setup.client_disconnect();
454KJ_TEST(
"Make simultaneous IPC calls on single remote thread")
458 std::promise<void> signal;
460 foo->initThreadMap();
463 setup.server->m_impl->m_fn = [&] {};
466 Thread::Client *callback_thread, *request_thread;
467 foo->m_context.loop->sync([&] {
468 Lock lock(tc.waiter->m_mutex);
469 callback_thread = &tc.callback_threads.at(foo->m_context.connection)->m_client;
470 request_thread = &tc.request_threads.at(foo->m_context.connection)->m_client;
474 std::atomic<int> expected = 100;
476 setup.server->m_impl->m_int_fn = [&](
int n) {
482 auto client{foo->m_client};
483 std::atomic<size_t> running{3};
484 foo->m_context.loop->sync([&]
486 for (
size_t i = 0; i < running; i++)
488 auto request{
client.callIntFnAsyncRequest()};
489 auto context{request.initContext()};
490 context.setCallbackThread(*callback_thread);
491 context.setThread(*request_thread);
492 request.setArg(100 * (i+1));
493 foo->m_context.loop->m_task_set->add(request.send().then(
494 [&running, &tc, i](
auto&& results) {
495 assert(results.getResult() == static_cast<int32_t>(100 * (i+1)));
497 Lock lock(tc.waiter->m_mutex);
498 tc.waiter->m_cv.notify_all();
503 Lock lock(tc.waiter->m_mutex);
504 tc.waiter->wait(lock, [&running] {
return running == 0; });
506 KJ_EXPECT(expected == 400);
509KJ_TEST(
"Call async IPC method dispatched to pool thread")
516 foo->initThreadMap();
517 setup.server->m_impl->m_int_fn = [](
int n) {
return n * 2; };
520 std::atomic<size_t> running{3};
521 std::promise<void> pool_ready;
522 foo->m_context.loop->sync([&] {
523 auto pool_req = foo->m_context.connection->m_thread_map.makePoolRequest();
524 pool_req.setCount(2);
525 foo->m_context.loop->m_task_set->add(
526 pool_req.send().then([&](
auto&&) { pool_ready.set_value(); }));
528 pool_ready.get_future().get();
532 auto client{foo->m_client};
533 foo->m_context.loop->sync([&] {
534 for (
size_t i = 0; i < running; ++i) {
535 auto request{
client.callIntFnAsyncRequest()};
536 request.initContext();
537 request.setArg(
static_cast<int32_t
>(i + 1));
538 foo->m_context.loop->m_task_set->add(request.send().then(
539 [&running, &tc, i](
auto&& results) {
540 assert(results.getResult() == static_cast<int32_t>((i + 1) * 2));
542 Lock lock(tc.waiter->m_mutex);
543 tc.waiter->m_cv.notify_all();
548 Lock lock(tc.waiter->m_mutex);
549 tc.waiter->wait(lock, [&running] {
return running == 0; });
553KJ_TEST(
"Call async IPC method without thread or pool errors correctly")
557 setup.server->m_impl->m_fn = [] {};
561 std::promise<void> done;
562 bool error_thrown{
false};
563 foo->m_context.loop->sync([&] {
564 auto request{foo->m_client.callFnAsyncRequest()};
565 request.initContext();
566 foo->m_context.loop->m_task_set->add(
568 [&](
auto&&) { done.set_value(); },
569 [&](kj::Exception&& e) {
571 KJ_EXPECT(std::string_view{e.getDescription().cStr()}.find(
572 "no thread specified and no pool configured") != std::string_view::npos);
576 done.get_future().get();
577 KJ_EXPECT(error_thrown);
Object holding network & rpc state associated with either an incoming server connection,...
Event loop implementation.
kj::AsyncIoContext m_io_context
Capnp IO context.
std::function< void()> testing_hook_makethread
Hook called when ProxyServer<ThreadMap>::makeThread() is called.
std::function< void()> testing_hook_makethread_created
Hook called on the worker thread inside makeThread(), after the thread context is set up and thread_c...
std::function< void()> testing_hook_async_request_done
Hook called on the worker thread just before returning results.
std::function< void()> testing_hook_async_request_start
Hook called on the worker thread when it starts to execute an async request.
Event loop smart pointer automatically managing m_num_refs.
Test setup class creating a two way connection between a ProxyServer<FooInterface> object and a Proxy...
std::promise< std::unique_ptr< ProxyClient< messages::FooInterface > > > client_promise
TestSetup(bool client_owns_connection=true)
std::function< void()> server_disconnect
std::function< void()> client_disconnect
std::thread thread
Thread variable should be after other struct members so the thread does not start until the other mem...
std::function< void()> server_disconnect_later
ProxyServer< messages::FooInterface > * server
std::unique_ptr< ProxyClient< messages::FooInterface > > client
#define EXPECT_EXCEPTION(call, message)
Assert that a call throws std::runtime_error with the given message.
std::unique_ptr< ProxyClient< messages::FooInterface > > client
std::promise< std::unique_ptr< ProxyClient< messages::FooInterface > > > client_promise
std::thread thread
Thread variable should be after other struct members so the thread does not start until the other mem...
std::vector< char > FooData
constexpr auto kMP_MAJOR_VERSION
Check version.h header values.
constexpr auto kMP_MINOR_VERSION
Functions to serialize / deserialize common bitcoin types.
thread_local ThreadContext g_thread_context
KJ_TEST("SpawnProcess does not run callback in child")
Log level
The severity level of this message.
std::string message
Message to be logged.
Mapping from capnp interface type to proxy client implementation (specializations are generated by pr...
Vat id for server side of connection.
The thread_local ThreadContext g_thread_context struct provides information about individual threads ...
std::map< std::string, int > map_string_int
std::unordered_set< int > unordered_set_int
std::optional< int > optional_int
std::vector< bool > vector_bool
Major and minor version numbers.
#define MP_MAJOR_VERSION
Major version number.
#define MP_MINOR_VERSION
Minor version number.