/* * unit_test.cpp — AntPeer C++17 unit tests (82 tests) * * Copyright (c) 2025-2026 Are Bjørby * SPDX-License-Identifier: MIT */ #include "antpeer.hpp" #include "Ed25519Utils.hpp" #include #include #include #include #include #include #include #include #include #include /* Reaching PAST antpeer, deliberately. test_membus_open() creates a named * POSIX shared-memory object and antpeer offers nothing that removes one — * Client::close() closes a session, not a bus. So a consumer that opens a bus * cannot close the loop through the API it opened with, and has to know the * layer underneath. That asymmetry is the root of this leak and is filed as * base/antpeer#ANTPEER-OPENS-A-NAMED-BUS-AND-CANNOT-DESTROY-ONE; this include * is the workaround, not the fix. */ static constexpr const char* TEST_BUS_NAME = "/antpeer_test"; static constexpr const char* TEST_OID = "test-oid"; static constexpr const char* TEST_DID = "test-did"; static constexpr const char* TEST_IID = "TS001"; static constexpr const char* TEST_CLI_IID = "CL001"; static constexpr const char* TEST_SERVICE = "TestService"; /* ── shm hygiene ────────────────────────────────────────────────────── * A POSIX shm object OUTLIVES the process that created it — that is the point * of naming one, and why membus's destructor does not unlink: closing a handle * must not destroy a rendezvous other peers may still be looking for. So a * suite that opens buses and exits leaves one kernel object per bus behind. * Measured 2026-08-19: a clean run took /dev/shm from 0 to 55 segments, 736 KB, * and the previous day's set was still there * (ANTPEER-TESTS-LEAK-A-SHM-SEGMENT-PER-BUS). * * THE NAMES ARE FIXED, which is what makes this more than untidiness: the next * run does not ignore a stale segment, it REUSES it, so a test can pass against * state its own earlier invocation left behind. * * REGISTERED BEFORE THE OPEN, not after. membus_open creates the object and can * still throw afterwards; registering first means a failed open is swept too. * Cleanup is idempotent — the remove half tolerates ENOENT — so a name that * several tests share costs nothing extra. * * IT SWEEPS THROUGH antpeer's OWN remove half since 3.3.0. Until then this * file had to `#include ` and call `membus::destroy`, * reaching past the library under test to undo what that library had done — * which is how ANTPEER-OPENS-A-NAMED-BUS-AND-CANNOT-DESTROY-ONE was found. * A suite that cannot clean up through the API it exercises is evidence about * that API, not a detail of the suite. * * THIS IS NOT THE GUARANTEE. It is a convenience that covers the one entry * point the suite uses today; a future test opening a bus another way leaks * again and this code cannot know. What catches that is scripts/shm-residue- * check.sh, which measures the directory rather than trusting the callers. */ static std::vector& opened_buses() { static std::vector v; return v; } static std::unique_ptr test_membus_open( std::string_view name, size_t size = antpeer::MEMBUS_DEFAULT_SIZE) { opened_buses().emplace_back(name); return antpeer::membus_open(name, size); } struct ShmSweeper { ~ShmSweeper() { for (const auto& n : opened_buses()) { /* A destructor may not throw, and a failure here must not mask a * test result: the residue check is what reports an unswept name. */ try { antpeer::membus_destroy(n); } catch (...) { } } } }; static int tests_run = 0; static int tests_passed = 0; static int tests_failed = 0; #define ASSERT(cond, msg) \ do { \ if (!(cond)) { \ std::fprintf(stderr, " FAIL: %s (line %d)\n", msg, __LINE__); \ return 0; \ } \ } while (0) #define RUN_TEST(fn) \ do { \ tests_run++; \ std::fprintf(stderr, " [%02d] %-50s ", tests_run, #fn); \ if (fn()) { tests_passed++; std::fprintf(stderr, "PASS\n"); } \ else { tests_failed++; std::fprintf(stderr, "FAILED\n"); }\ } while (0) /* ═══════════════════════════════════════════════════════════════════════ * Bus tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_bus_membus_open_close() { auto bus = test_membus_open(TEST_BUS_NAME, antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "membus open returned nullptr"); ASSERT(bus->is_open(), "membus not open after open"); ASSERT(!bus->name().empty(), "bus name is empty"); ASSERT(bus->name() == TEST_BUS_NAME, "bus name mismatch"); return 1; } static int test_bus_write_read() { auto bus = test_membus_open(TEST_BUS_NAME, antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "membus open failed"); const char* msg = "HELLO"; ssize_t w = bus->write(msg, std::strlen(msg)); ASSERT(w > 0, "write failed"); char buf[256] = {}; ssize_t r = bus->read(buf, sizeof(buf)); ASSERT(r > 0, "read returned no data"); ASSERT(std::memcmp(buf, msg, std::strlen(msg)) == 0, "read data mismatch"); return 1; } static int test_bus_reopen() { auto bus = test_membus_open(TEST_BUS_NAME, antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "membus open failed"); ASSERT(bus->is_open(), "not open initially"); bool ok = bus->reopen(); ASSERT(ok, "reopen failed"); ASSERT(bus->is_open(), "not open after reopen"); return 1; } static int test_bus_read_empty() { auto bus = test_membus_open("/antpeer_test_empty", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "membus open failed"); char buf[256]; ssize_t r = bus->read(buf, sizeof(buf)); ASSERT(r == 0, "expected 0 from empty bus"); return 1; } static int test_bus_empty_name() { try { auto bus = test_membus_open("", antpeer::MEMBUS_DEFAULT_SIZE); return 0; /* Should have thrown */ } catch (const std::system_error& e) { return (e.code().value() == EINVAL); } } static int test_bus_zero_size() { auto bus = test_membus_open("/antpeer_test_zero", 0); ASSERT(bus != nullptr, "membus open with size 0 should use default"); ASSERT(bus->is_open(), "bus not open"); return 1; } static int test_bus_factory_throws() { /* Empty name for sockbus should also throw */ try { auto bus = antpeer::sockbus_open(""); return 0; } catch (const std::system_error& e) { return (e.code().value() == EINVAL); } } static int test_bus_large_write_read() { auto bus = test_membus_open("/antpeer_test_large", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); char data[4096]; std::memset(data, 'A', sizeof(data)); ssize_t w = bus->write(data, sizeof(data)); ASSERT(w > 0, "large write failed"); char buf[8192] = {}; ssize_t r = bus->read(buf, sizeof(buf)); ASSERT(r > 0, "large read returned no data"); ASSERT(r >= static_cast(sizeof(data)), "read less than written"); ASSERT(std::memcmp(buf, data, sizeof(data)) == 0, "large data mismatch"); return 1; } static int test_bus_multiple_writes() { auto bus = test_membus_open("/antpeer_test_multi_w", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); ASSERT(bus->write("AAA", 3) > 0, "write 1 failed"); ASSERT(bus->write("BBB", 3) > 0, "write 2 failed"); char buf[256] = {}; ssize_t total = 0; ssize_t r; while ((r = bus->read(buf + total, sizeof(buf) - static_cast(total))) > 0) total += r; ASSERT(total >= 6, "didn't read all written data"); return 1; } static int test_bus_reopen_preserves_name() { auto bus = test_membus_open("/antpeer_test_rn", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); auto name_before = bus->name(); ASSERT(!name_before.empty(), "name before reopen is empty"); ASSERT(bus->reopen(), "reopen failed"); auto name_after = bus->name(); ASSERT(!name_after.empty(), "name after reopen is empty"); ASSERT(name_before == name_after, "name changed after reopen"); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Server tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_server_create_destroy() { antpeer::Server s; (void)s; return 1; } static int test_server_init() { auto bus = test_membus_open("/antpeer_test_srv", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); const char* bid = s.bid(); ASSERT(bid != nullptr, "BID is NULL after init"); ASSERT(std::strlen(bid) > 0, "BID is empty"); return 1; } static int test_server_establish() { auto bus = test_membus_open("/antpeer_test_est", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); ASSERT(s.bid() != nullptr, "BID should be assigned after establish"); return 1; } static int test_server_stop() { auto bus = test_membus_open("/antpeer_test_stop", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); std::thread t([&s]() { s.run(); }); usleep(50000); s.stop(); t.join(); return 1; } static int test_server_broadcast() { auto bus = test_membus_open("/antpeer_test_bc", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "BcSvc"); ASSERT(rc == 0, "server init failed"); rc = s.broadcast("HELLO_ALL"); ASSERT(rc == 0, "broadcast failed"); return 1; } struct HandlerCtx { std::string reply_body; int call_count = 0; }; static int test_server_handler_reply() { auto sbus = test_membus_open("/antpeer_test_hr", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_hr", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; HandlerCtx hctx{"PONG", 0}; s.on_request([&hctx](antpeer::Server& srv, std::string_view sid, std::string_view) { hctx.call_count++; srv.reply(sid, hctx.reply_body); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "EchoSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("EchoSvc", 3000); ASSERT(rc == 0, "discover failed"); ASSERT(c.connected(), "not connected after discover"); char reply[256] = {}; int n = c.call("PING", 3000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); ASSERT(std::strcmp(reply, "PONG") == 0, "reply mismatch"); ASSERT(hctx.call_count == 1, "handler not called exactly once"); s.stop(); t.join(); return 1; } static int test_server_session_handler() { auto sbus = test_membus_open("/antpeer_test_sh", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_sh", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; int session_calls = 0; s.on_session([&session_calls](antpeer::Server& srv, std::string_view sid, uint32_t, std::string_view) { session_calls++; srv.reply(sid, "SESSION_OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "SessSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("SessSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("REQ1", 3000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); ASSERT(std::strcmp(reply, "SESSION_OK") == 0, "reply mismatch"); ASSERT(session_calls >= 1, "session handler not called"); s.stop(); t.join(); return 1; } static int test_server_default_state() { antpeer::Server s; ASSERT(s.bid() == nullptr, "bid should be NULL before init"); ASSERT(s.reply("x", "y") == -1, "reply before init should return -1"); ASSERT(s.notify("x", "y") == -1, "notify before init should return -1"); ASSERT(s.broadcast("x") == -1, "broadcast before init should return -1"); ASSERT(s.finish("x") == -1, "finish before init should return -1"); ASSERT(s.offer_additional("x") == -1, "offer before init should return -1"); return 1; } static int test_server_init_bad_args() { auto bus = test_membus_open("/antpeer_test_bad", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; ASSERT(s.init(*bus, "", TEST_DID, TEST_IID, TEST_SERVICE) == -1, "init with empty oid should fail"); ASSERT(s.init(*bus, TEST_OID, "", TEST_IID, TEST_SERVICE) == -1, "init with empty did should fail"); ASSERT(s.init(*bus, TEST_OID, TEST_DID, "", TEST_SERVICE) == -1, "init with empty iid should fail"); ASSERT(s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "") == -1, "init with empty service should fail"); return 1; } static int test_server_double_destroy() { { antpeer::Server s; } { antpeer::Server s2; } return 1; } static int test_server_reply_unknown_sid() { auto bus = test_membus_open("/antpeer_test_runk", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); ASSERT(s.reply("UNKNOWN_SID", "body") == -1, "reply with unknown SID should fail"); return 1; } static int test_server_finish_unknown_sid() { auto bus = test_membus_open("/antpeer_test_funk", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); ASSERT(s.finish("UNKNOWN_SID") == -1, "finish with unknown SID should fail"); return 1; } static int test_server_offer_additional() { auto bus = test_membus_open("/antpeer_test_oa", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "MainSvc"); ASSERT(rc == 0, "init failed"); rc = s.offer_additional("ExtraSvc"); ASSERT(rc == 0, "offer_additional failed"); ASSERT(s.offer_additional("") == -1, "offer with empty description should fail"); return 1; } static int test_server_bid_before_init() { antpeer::Server s; const char* bid = s.bid(); ASSERT(bid == nullptr, "BID should be NULL before init"); return 1; } static int test_server_set_expiry() { antpeer::Server s; s.set_session_expiry(120); s.set_session_expiry(0); s.set_session_expiry(3600); return 1; } static int test_server_notify_unknown_sid() { auto bus = test_membus_open("/antpeer_test_nunk", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, TEST_SERVICE); ASSERT(rc == 0, "init failed"); ASSERT(s.notify("UNKNOWN_SID", "event") == -1, "notify with unknown SID should fail"); return 1; } static int test_server_multiple_broadcasts() { auto bus = test_membus_open("/antpeer_test_mbc", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "MbcSvc"); ASSERT(rc == 0, "init failed"); for (int i = 0; i < 5; i++) { char event[32]; std::snprintf(event, sizeof(event), "EVENT_%d", i); rc = s.broadcast(event); ASSERT(rc == 0, "broadcast failed"); } return 1; } static int test_server_broadcast_empty_event() { auto bus = test_membus_open("/antpeer_test_bcn", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Server s; int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "BcnSvc"); ASSERT(rc == 0, "init failed"); rc = s.broadcast(""); ASSERT(rc == 0, "broadcast with empty event should succeed"); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Client tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_client_create_destroy() { antpeer::Client c; ASSERT(c.state() == antpeer::SessionState::Disconnected, "initial state not Disconnected"); return 1; } static int test_client_init() { auto bus = test_membus_open("/antpeer_test_cli", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; int rc = c.init(*bus, TEST_OID, TEST_DID, TEST_CLI_IID); ASSERT(rc == 0, "init failed"); return 1; } static int test_client_discover() { auto sbus = test_membus_open("/antpeer_test_disc", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_disc", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "DiscSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("DiscSvc", 3000); ASSERT(rc == 0, "discover failed"); ASSERT(c.connected(), "not connected"); const char* pbid = c.peer_bid(); ASSERT(pbid != nullptr, "peer BID is NULL"); ASSERT(std::strlen(pbid) > 0, "peer BID is empty"); s.stop(); t.join(); return 1; } static int test_client_call_timeout() { auto sbus = test_membus_open("/antpeer_test_to", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_to", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TimeoutSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("TimeoutSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("PING", 500, reply, sizeof(reply)); ASSERT(n == 0, "expected timeout return (0)"); s.stop(); t.join(); return 1; } static int test_client_fire() { auto sbus = test_membus_open("/antpeer_test_fire", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_fire", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; HandlerCtx hctx{"ACK", 0}; s.on_request([&hctx](antpeer::Server& srv, std::string_view sid, std::string_view) { hctx.call_count++; srv.reply(sid, hctx.reply_body); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "FireSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("FireSvc", 3000); ASSERT(rc == 0, "discover failed"); rc = c.fire("FIRE_CMD"); ASSERT(rc == 0, "fire failed"); usleep(100000); ASSERT(hctx.call_count >= 1, "handler not called after fire"); s.stop(); t.join(); return 1; } static int test_client_close_reopen() { auto bus = test_membus_open("/antpeer_test_cr", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; int rc = c.init(*bus, TEST_OID, TEST_DID, TEST_CLI_IID); ASSERT(rc == 0, "init failed"); ASSERT(!c.connected(), "should not be connected without discover"); c.close(); ASSERT(c.state() == antpeer::SessionState::Disconnected, "not disconnected after close"); ASSERT(c.peer_bid() == nullptr, "peer BID not NULL after close"); return 1; } static int test_client_notify() { auto sbus = test_membus_open("/antpeer_test_ntf", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_ntf", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "NtfSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; int notify_count = 0; c.on_notify([¬ify_count](std::string_view, std::string_view) { notify_count++; }); rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("NtfSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("HELLO", 3000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); ASSERT(std::strcmp(reply, "HELLO") == 0, "echo mismatch"); ASSERT(notify_count == 0, "notify called when no N-verb was sent"); s.stop(); t.join(); return 1; } static int test_client_multi_call() { auto sbus = test_membus_open("/antpeer_test_mc", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_mc", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "EchoSvc2"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("EchoSvc2", 3000); ASSERT(rc == 0, "discover failed"); char reply[256]; for (int i = 0; i < 5; i++) { char cmd[32]; std::snprintf(cmd, sizeof(cmd), "MSG_%d", i); std::memset(reply, 0, sizeof(reply)); int n = c.call(cmd, 3000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); ASSERT(std::strcmp(reply, cmd) == 0, "echo mismatch"); } s.stop(); t.join(); return 1; } static std::atomic shutdown_called{0}; static int test_client_required_shutdown() { auto sbus = test_membus_open("/antpeer_test_req", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_req", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "ReqSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; shutdown_called.store(0); c.set_required(true); c.on_shutdown([]() { shutdown_called.store(1); }); rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("ReqSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("WILL_TIMEOUT", 500, reply, sizeof(reply)); ASSERT(n == -1, "expected -1 for required peer failure"); ASSERT(shutdown_called.load() == 1, "shutdown callback not called"); s.stop(); t.join(); return 1; } static int test_client_default_state() { antpeer::Client c; ASSERT(!c.connected(), "connected() should be false"); ASSERT(c.state() == antpeer::SessionState::Disconnected, "state should be Disconnected"); ASSERT(c.peer_bid() == nullptr, "peer_bid should be NULL"); ASSERT(c.ensure_connected(100) == -1, "ensure_connected should return -1"); return 1; } static int test_client_init_bad_args() { auto bus = test_membus_open("/antpeer_test_cbad", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; ASSERT(c.init(*bus, "", TEST_DID, TEST_CLI_IID) == -1, "init with empty oid should fail"); ASSERT(c.init(*bus, TEST_OID, "", TEST_CLI_IID) == -1, "init with empty did should fail"); ASSERT(c.init(*bus, TEST_OID, TEST_DID, "") == -1, "init with empty iid should fail"); return 1; } static int test_client_double_destroy() { { antpeer::Client c; } { antpeer::Client c2; } return 1; } static int test_client_call_not_connected() { antpeer::Client c; char reply[64] = {}; int n = c.call("test", 100, reply, sizeof(reply)); ASSERT(n <= 0, "call without connection should not succeed"); return 1; } static int test_client_fire_not_connected() { antpeer::Client c; int rc = c.fire("test"); ASSERT(rc == -1, "fire without connection should fail"); return 1; } static int test_client_peer_bid_before_discover() { auto bus = test_membus_open("/antpeer_test_pbid", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; int rc = c.init(*bus, TEST_OID, TEST_DID, TEST_CLI_IID); ASSERT(rc == 0, "init failed"); ASSERT(c.peer_bid() == nullptr, "peer BID should be NULL before discover"); ASSERT(!c.connected(), "should not be connected before discover"); return 1; } static int test_client_discover_timeout() { auto bus = test_membus_open("/antpeer_test_dto", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; int rc = c.init(*bus, TEST_OID, TEST_DID, TEST_CLI_IID); ASSERT(rc == 0, "init failed"); rc = c.discover("NonExistentSvc", 500); ASSERT(rc == -1, "discover with no server should timeout/fail"); ASSERT(!c.connected(), "should not be connected after failed discover"); return 1; } static int test_client_discover_empty_service() { auto bus = test_membus_open("/antpeer_test_dns", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus != nullptr, "bus open failed"); antpeer::Client c; int rc = c.init(*bus, TEST_OID, TEST_DID, TEST_CLI_IID); ASSERT(rc == 0, "init failed"); rc = c.discover("", 500); ASSERT(rc == -1, "discover with empty service should fail"); return 1; } static int test_client_state_transitions() { auto sbus = test_membus_open("/antpeer_test_st", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_st", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "StateSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; ASSERT(c.state() == antpeer::SessionState::Disconnected, "initial state not Disconnected"); rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "init failed"); rc = c.discover("StateSvc", 3000); ASSERT(rc == 0, "discover failed"); ASSERT(c.state() == antpeer::SessionState::Connected, "state not Connected after discover"); c.close(); ASSERT(c.state() == antpeer::SessionState::Disconnected, "state not Disconnected after close"); s.stop(); t.join(); return 1; } static int test_client_call_null_command() { antpeer::Client c; int n = c.call({nullptr, 0}, 100, nullptr, 0); ASSERT(n <= 0, "call with null command should not succeed"); return 1; } static int test_client_fire_null_command() { antpeer::Client c; int rc = c.fire({nullptr, 0}); ASSERT(rc == -1, "fire with null command should fail"); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Integration depth tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_server_finish_handler_called() { auto sbus = test_membus_open("/antpeer_test_fh", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_fh", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; HandlerCtx hctx{"OK", 0}; s.on_request([&hctx](antpeer::Server& srv, std::string_view sid, std::string_view) { hctx.call_count++; srv.reply(sid, hctx.reply_body); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "FinSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("FinSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("CMD", 3000, reply, sizeof(reply)); ASSERT(n > 0, "call failed"); ASSERT(std::strcmp(reply, "OK") == 0, "reply mismatch"); c.close(); usleep(50000); s.stop(); t.join(); return 1; } static int test_multi_client_same_server() { auto sbus = test_membus_open("/antpeer_test_2cl", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus1 = test_membus_open("/antpeer_test_2cl", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus2 = test_membus_open("/antpeer_test_2cl", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus1 && cbus2, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "MultiSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c1, c2; rc = c1.init(*cbus1, "cli1-oid", "cli1-did", "C1001"); ASSERT(rc == 0, "client 1 init failed"); rc = c1.discover("MultiSvc", 3000); ASSERT(rc == 0, "client 1 discover failed"); rc = c2.init(*cbus2, "cli2-oid", "cli2-did", "C2001"); ASSERT(rc == 0, "client 2 init failed"); rc = c2.discover("MultiSvc", 3000); ASSERT(rc == 0, "client 2 discover failed"); char reply1[256] = {}, reply2[256] = {}; int n1 = c1.call("FROM_C1", 3000, reply1, sizeof(reply1)); int n2 = c2.call("FROM_C2", 3000, reply2, sizeof(reply2)); ASSERT(n1 > 0, "client 1 call failed"); ASSERT(n2 > 0, "client 2 call failed"); ASSERT(std::strcmp(reply1, "FROM_C1") == 0, "client 1 echo mismatch"); ASSERT(std::strcmp(reply2, "FROM_C2") == 0, "client 2 echo mismatch"); s.stop(); t.join(); return 1; } static int test_call_large_payload() { auto sbus = test_membus_open("/antpeer_test_lp", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_lp", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "LargeSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("LargeSvc", 3000); ASSERT(rc == 0, "discover failed"); char big_cmd[1024]; std::memset(big_cmd, 'X', sizeof(big_cmd) - 1); big_cmd[sizeof(big_cmd) - 1] = '\0'; char reply[2048] = {}; int n = c.call(big_cmd, 3000, reply, sizeof(reply)); ASSERT(n > 0, "large payload call failed"); ASSERT(std::strcmp(reply, big_cmd) == 0, "large payload echo mismatch"); s.stop(); t.join(); return 1; } static int test_call_empty_command() { auto sbus = test_membus_open("/antpeer_test_ec", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_ec", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { srv.reply(sid, cmd); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "EmptySvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("EmptySvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("", 3000, reply, sizeof(reply)); ASSERT(n >= 0, "empty command call should not error"); s.stop(); t.join(); return 1; } static int test_client_call_small_reply_buffer() { auto sbus = test_membus_open("/antpeer_test_srb", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_srb", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "LONG_RESPONSE_BODY"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TruncSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("TruncSvc", 3000); ASSERT(rc == 0, "discover failed"); char reply[8] = {}; int n = c.call("CMD", 3000, reply, sizeof(reply)); ASSERT(n > 0, "call failed with small buffer"); ASSERT(n <= 7, "should not write more than buffer allows"); ASSERT(reply[7] == '\0', "reply not null-terminated"); s.stop(); t.join(); return 1; } static int test_client_call_null_reply_buffer() { auto sbus = test_membus_open("/antpeer_test_nrb", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_nrb", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; HandlerCtx hctx{"REPLY", 0}; s.on_request([&hctx](antpeer::Server& srv, std::string_view sid, std::string_view) { hctx.call_count++; srv.reply(sid, hctx.reply_body); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "NullBufSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("NullBufSvc", 3000); ASSERT(rc == 0, "discover failed"); int n = c.call("CMD", 3000, nullptr, 0); ASSERT(n == 0, "call with NULL buffer should return 0 bytes written"); ASSERT(hctx.call_count == 1, "handler not called"); s.stop(); t.join(); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * IID / Misc tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_iid_caller_provided() { auto sbus = test_membus_open("/antpeer_test_iidp", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_iidp", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; int rc = s.init(*sbus, TEST_OID, TEST_DID, "MY999", "IidSvc"); ASSERT(rc == 0, "server init with custom IID failed"); ASSERT(s.bid() != nullptr, "server BID is NULL"); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", "XY123"); ASSERT(rc == 0, "client init with custom IID failed"); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Identity tests (Verify protocol) * ═══════════════════════════════════════════════════════════════════════ */ static int test_peer_identity_struct() { antpeer::PeerIdentity pi; ASSERT(!pi.verified(), "default PeerIdentity should not be verified"); ASSERT(pi.bid.empty(), "default bid should be empty"); pi.oid = "TTEi"; pi.did = "MAIN"; pi.iid = "SVC00"; pi.bid = "A2"; ASSERT(pi.verified(), "should be verified when oid is set"); return 1; } static int test_accept_callback() { auto sbus = test_membus_open("/antpeer_test_acb", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_acb", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; std::string accepted_bid; s.on_accept([&accepted_bid](std::string_view peer_bid) { accepted_bid = std::string(peer_bid); }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "AcbSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "test-oid", "test-did", "TC001"); ASSERT(rc == 0, "client init failed"); rc = c.discover("AcbSvc", 3000); ASSERT(rc == 0, "discover failed"); /* Client Accept sends sender BID — give server time to process */ char reply[64] = {}; c.call("PING", 3000, reply, sizeof(reply)); s.stop(); t.join(); ASSERT(!accepted_bid.empty(), "accept callback not fired"); return 1; } static int test_verify_roundtrip() { auto sbus = test_membus_open("/antpeer_test_vrt", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_vrt", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; antpeer::PeerIdentity verified_id; s.on_verify([&verified_id](const antpeer::PeerIdentity& id) { verified_id = id; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "VrtSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "PEER-ORG", "PEER-DEV", "SN042"); ASSERT(rc == 0, "client init failed"); rc = c.discover("VrtSvc", 3000); ASSERT(rc == 0, "discover failed"); /* call() triggers Accept → auto-Verify → response → callback */ char reply[64] = {}; c.call("PING", 3000, reply, sizeof(reply)); s.stop(); t.join(); ASSERT(verified_id.verified(), "verify callback not fired"); ASSERT(verified_id.oid == "PEER-ORG", "oid mismatch"); ASSERT(verified_id.did == "PEER-DEV", "did mismatch"); ASSERT(verified_id.iid == "SN042", "iid mismatch"); ASSERT(!verified_id.bid.empty(), "bid should be set"); /* peer_identity lookup */ auto* pi = s.peer_identity(verified_id.bid); ASSERT(pi != nullptr, "peer_identity returned nullptr"); ASSERT(pi->oid == "PEER-ORG", "cached oid mismatch"); return 1; } static int test_server_responds_to_verify() { /* Two servers on same bus — one verifies the other */ auto bus1 = test_membus_open("/antpeer_test_srv2v", antpeer::MEMBUS_DEFAULT_SIZE); auto bus2 = test_membus_open("/antpeer_test_srv2v", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus1 && bus2, "bus open failed"); antpeer::Server s1; s1.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s1.init(*bus1, "ORG-A", "DEV-A", "S001", "Svc1"); ASSERT(rc == 0, "s1 init failed"); antpeer::Server s2; antpeer::PeerIdentity s2_verified; s2.on_verify([&s2_verified](const antpeer::PeerIdentity& id) { s2_verified = id; }); s2.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); rc = s2.init(*bus2, "ORG-B", "DEV-B", "S002", "Svc2"); ASSERT(rc == 0, "s2 init failed"); /* Run both servers briefly so they can exchange Verify */ std::thread t1([&s1]() { s1.run(); }); std::thread t2([&s2]() { s2.run(); }); /* s2 explicitly verifies s1 */ rc = s2.send_verify(s1.bid()); ASSERT(rc == 0, "send_verify failed"); usleep(200000); /* let Verify round-trip complete */ s1.stop(); s2.stop(); t1.join(); t2.join(); ASSERT(s2_verified.verified(), "s2 did not receive verify response"); ASSERT(s2_verified.oid == "ORG-A", "s2 verified oid mismatch"); ASSERT(s2_verified.did == "DEV-A", "s2 verified did mismatch"); ASSERT(s2_verified.iid == "S001", "s2 verified iid mismatch"); return 1; } static int test_unverified_peer_no_identity() { antpeer::Server s; auto bus = test_membus_open("/antpeer_test_unv", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(bus, "bus open failed"); int rc = s.init(*bus, TEST_OID, TEST_DID, TEST_IID, "UnvSvc"); ASSERT(rc == 0, "server init failed"); auto* pi = s.peer_identity("NONEXISTENT"); ASSERT(pi == nullptr, "should return nullptr for unknown BID"); auto pbid = s.peer_bid_for_session("UNKNOWN_SID"); ASSERT(pbid.empty(), "should return empty for unknown SID"); return 1; } static int test_sid_peer_resolution() { auto sbus = test_membus_open("/antpeer_test_spr", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_spr", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; std::string resolved_bid; std::string call_sid; s.on_verify([&resolved_bid](const antpeer::PeerIdentity& id) { resolved_bid = id.bid; }); s.on_session([&call_sid](antpeer::Server& srv, std::string_view sid, uint32_t, std::string_view) { call_sid = std::string(sid); srv.reply(sid, "OK"); srv.finish(sid); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "SprSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "TEST-ORG", "TEST-DEV", "SN100"); ASSERT(rc == 0, "client init failed"); rc = c.discover("SprSvc", 3000); ASSERT(rc == 0, "discover failed"); auto resp = c.call("TEST", 3000); ASSERT(!resp.empty(), "call failed"); s.stop(); t.join(); ASSERT(!resolved_bid.empty(), "verify not received"); ASSERT(!call_sid.empty(), "session not received"); /* The SID from the call should resolve to the verified peer BID */ auto peer = s.peer_bid_for_session(call_sid); ASSERT(!peer.empty(), "SID not resolved to peer"); ASSERT(peer == resolved_bid, "SID resolved to wrong peer"); return 1; } static int test_multiple_peers_verify() { auto sbus = test_membus_open("/antpeer_test_mpv", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus1 = test_membus_open("/antpeer_test_mpv", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus2 = test_membus_open("/antpeer_test_mpv", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus1 && cbus2, "bus open failed"); antpeer::Server s; std::vector verified_peers; s.on_verify([&verified_peers](const antpeer::PeerIdentity& id) { verified_peers.push_back(id); }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "MpvSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c1; rc = c1.init(*cbus1, "ORG-1", "DEV-1", "SN001"); ASSERT(rc == 0, "c1 init failed"); rc = c1.discover("MpvSvc", 3000); ASSERT(rc == 0, "c1 discover failed"); char reply[64] = {}; c1.call("P1", 3000, reply, sizeof(reply)); antpeer::Client c2; rc = c2.init(*cbus2, "ORG-2", "DEV-2", "SN002"); ASSERT(rc == 0, "c2 init failed"); rc = c2.discover("MpvSvc", 3000); ASSERT(rc == 0, "c2 discover failed"); c2.call("P2", 3000, reply, sizeof(reply)); s.stop(); t.join(); ASSERT(verified_peers.size() >= 2, "should have verified 2 peers"); /* Both peers should be in the identity cache via their BIDs */ auto* pi1 = s.peer_identity(verified_peers[0].bid); auto* pi2 = s.peer_identity(verified_peers[1].bid); ASSERT(pi1 != nullptr, "peer 1 not in cache"); ASSERT(pi2 != nullptr, "peer 2 not in cache"); /* Order may vary — check that both ORGs are present */ bool has_org1 = (pi1->oid == "ORG-1") || (pi2->oid == "ORG-1"); bool has_org2 = (pi1->oid == "ORG-2") || (pi2->oid == "ORG-2"); ASSERT(has_org1, "ORG-1 not found in cache"); ASSERT(has_org2, "ORG-2 not found in cache"); return 1; } static int test_send_verify_explicit() { auto sbus = test_membus_open("/antpeer_test_sve", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_sve", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; antpeer::PeerIdentity manual_id; s.on_verify([&manual_id](const antpeer::PeerIdentity& id) { manual_id = id; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "SveSvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "MANUAL-ORG", "MANUAL-DEV", "MN001"); ASSERT(rc == 0, "client init failed"); rc = c.discover("SveSvc", 3000); ASSERT(rc == 0, "discover failed"); /* After discover+accept, auto-verify already fires. * This tests that explicit send_verify also works (second verify). */ char reply[64] = {}; c.call("PING", 3000, reply, sizeof(reply)); /* The auto-verify should have already populated manual_id */ ASSERT(manual_id.verified(), "auto-verify should have fired"); ASSERT(manual_id.oid == "MANUAL-ORG", "oid mismatch"); s.stop(); t.join(); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * FrameExtractor tests * ═══════════════════════════════════════════════════════════════════════ */ /* Helper: build a BLOB word for declaring tail size. * Wire: [SOW] * [SOR]D [SOU]B [SOB] [EOW] */ static std::vector make_blob_word(size_t tail_size) { std::string digits = std::to_string(tail_size); std::vector w; w.push_back(0x12); /* SOW */ w.push_back('*'); /* BLOB type */ w.push_back(0x04); /* SOR */ w.push_back('D'); /* Decimal radix */ w.push_back(0x07); /* SOU */ w.push_back('B'); /* Byte unit */ w.push_back(0x1A); /* SOB */ for (char c : digits) w.push_back(static_cast(c)); w.push_back(0x10); /* EOW */ return w; } /* Helper: build a complete Antheos message with optional tail. * [SOM] words... [EOM] tail_bytes */ static std::vector make_message( const std::vector>& words, const uint8_t* tail = nullptr, size_t tail_len = 0) { std::vector msg; msg.push_back(0x02); /* SOM */ for (auto& w : words) msg.insert(msg.end(), w.begin(), w.end()); msg.push_back(0x03); /* EOM */ if (tail && tail_len > 0) msg.insert(msg.end(), tail, tail + tail_len); return msg; } static int test_fe_create_destroy() { antpeer::FrameExtractor fe; ASSERT(fe.total_frames() == 0, "initial frame count not 0"); ASSERT(fe.parse_errors() == 0, "initial error count not 0"); return 1; } static int test_fe_simple_frame() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "ping frame build failed"); fe.feed(frame->data(), frame->size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == frame->size(), "frame size mismatch"); ASSERT(std::memcmp(captured.data(), frame->data(), frame->size()) == 0, "frame content mismatch"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_frame_with_tail() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); uint8_t tail_data[] = {0xAA, 0xBB, 0xCC, 0xDD, 0xEE}; auto blob = make_blob_word(5); auto msg = make_message({blob}, tail_data, 5); fe.feed(msg.data(), msg.size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == msg.size(), "frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "frame content mismatch"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_tail_with_som_eom_bytes() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); /* Tail contains SOM (0x02) and EOM (0x03) — must not confuse extractor */ uint8_t tail_data[] = {0x02, 0x03, 0x02, 0x03, 0xFF}; auto blob = make_blob_word(5); auto msg = make_message({blob}, tail_data, 5); fe.feed(msg.data(), msg.size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == msg.size(), "frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "tail with SOM/EOM bytes corrupted"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_multiple_frames() { antpeer::FrameExtractor fe; std::vector> frames; fe.on_frame([&](const uint8_t* data, size_t len) { frames.emplace_back(data, data + len); }); auto f1 = antheos::bus::ping("AAAA"); auto f2 = antheos::bus::ping("BBBB"); ASSERT(f1.has_value() && f2.has_value(), "frame build failed"); /* Feed both frames in a single buffer */ std::vector combined; combined.insert(combined.end(), f1->begin(), f1->end()); combined.insert(combined.end(), f2->begin(), f2->end()); fe.feed(combined.data(), combined.size()); ASSERT(frames.size() == 2, "expected 2 frames"); ASSERT(frames[0].size() == f1->size(), "frame 1 size mismatch"); ASSERT(frames[1].size() == f2->size(), "frame 2 size mismatch"); ASSERT(fe.total_frames() == 2, "frame count not 2"); return 1; } static int test_fe_partial_feed() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "frame build failed"); /* Feed one byte at a time */ for (size_t i = 0; i < frame->size(); i++) fe.feed(frame->data() + i, 1); ASSERT(!captured.empty(), "no frame emitted from byte-at-a-time feed"); ASSERT(captured.size() == frame->size(), "frame size mismatch"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_partial_feed_with_tail() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); uint8_t tail_data[] = {0x02, 0x03, 0x42}; auto blob = make_blob_word(3); auto msg = make_message({blob}, tail_data, 3); /* Feed one byte at a time */ for (size_t i = 0; i < msg.size(); i++) fe.feed(msg.data() + i, 1); ASSERT(!captured.empty(), "no frame emitted from byte-at-a-time feed"); ASSERT(captured.size() == msg.size(), "frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "content mismatch"); return 1; } static int test_fe_garbage_before_som() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "frame build failed"); /* Prepend garbage bytes */ std::vector buf = {0xFF, 0xFE, 0xFD, 0x41, 0x42}; buf.insert(buf.end(), frame->begin(), frame->end()); fe.feed(buf.data(), buf.size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == frame->size(), "garbage not skipped"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_large_tail() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); /* 8KB tail — exceeds antheos::Parser's internal 4KB buffer */ constexpr size_t TAIL_SIZE = 8192; std::vector tail_data(TAIL_SIZE); for (size_t i = 0; i < TAIL_SIZE; i++) tail_data[i] = static_cast(i & 0xFF); auto blob = make_blob_word(TAIL_SIZE); auto msg = make_message({blob}, tail_data.data(), TAIL_SIZE); fe.feed(msg.data(), msg.size()); ASSERT(!captured.empty(), "no frame emitted for large tail"); ASSERT(captured.size() == msg.size(), "large tail frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "large tail content mismatch"); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_multiple_blob_words() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); /* Two BLOB words: 3 bytes + 4 bytes = 7 bytes total tail */ auto blob1 = make_blob_word(3); auto blob2 = make_blob_word(4); uint8_t tail_data[] = {0x11, 0x22, 0x33, 0x44, 0x55, 0x66, 0x77}; auto msg = make_message({blob1, blob2}, tail_data, 7); fe.feed(msg.data(), msg.size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == msg.size(), "multi-blob frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "multi-blob content mismatch"); return 1; } static int test_fe_reset() { antpeer::FrameExtractor fe; size_t frame_count = 0; fe.on_frame([&](const uint8_t*, size_t) { frame_count++; }); auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "frame build failed"); /* Feed half a frame, then reset, then feed a complete frame */ size_t half = frame->size() / 2; fe.feed(frame->data(), half); fe.reset(); ASSERT(frame_count == 0, "frame emitted from partial before reset"); fe.feed(frame->data(), frame->size()); ASSERT(frame_count == 1, "frame not emitted after reset + full feed"); return 1; } static int test_fe_no_callback() { antpeer::FrameExtractor fe; /* No on_frame callback — should not crash */ auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "frame build failed"); fe.feed(frame->data(), frame->size()); ASSERT(fe.total_frames() == 1, "frame count not 1"); return 1; } static int test_fe_tail_then_next_frame() { antpeer::FrameExtractor fe; std::vector> frames; fe.on_frame([&](const uint8_t* data, size_t len) { frames.emplace_back(data, data + len); }); /* Frame 1: message with 3-byte tail */ uint8_t tail_data[] = {0x02, 0x03, 0x42}; auto blob = make_blob_word(3); auto msg1 = make_message({blob}, tail_data, 3); /* Frame 2: simple ping (no tail) */ auto msg2 = antheos::bus::ping("CCCC"); ASSERT(msg2.has_value(), "ping build failed"); /* Feed both concatenated */ std::vector combined; combined.insert(combined.end(), msg1.begin(), msg1.end()); combined.insert(combined.end(), msg2->begin(), msg2->end()); fe.feed(combined.data(), combined.size()); ASSERT(frames.size() == 2, "expected 2 frames"); ASSERT(frames[0].size() == msg1.size(), "frame 1 size mismatch"); ASSERT(frames[1].size() == msg2->size(), "frame 2 size mismatch"); return 1; } static int test_fe_hex_blob() { antpeer::FrameExtractor fe; std::vector captured; fe.on_frame([&](const uint8_t* data, size_t len) { captured.assign(data, data + len); }); /* BLOB word with hex radix: 0x0A = 10 bytes */ std::vector blob_hex; blob_hex.push_back(0x12); /* SOW */ blob_hex.push_back('*'); /* BLOB */ blob_hex.push_back(0x04); /* SOR */ blob_hex.push_back('H'); /* Hex radix */ blob_hex.push_back(0x07); /* SOU */ blob_hex.push_back('W'); /* Word unit */ blob_hex.push_back(0x1A); /* SOB */ blob_hex.push_back('A'); /* hex "A" = 10 */ blob_hex.push_back(0x10); /* EOW */ uint8_t tail_data[10]; std::memset(tail_data, 0x42, 10); auto msg = make_message({blob_hex}, tail_data, 10); fe.feed(msg.data(), msg.size()); ASSERT(!captured.empty(), "no frame emitted"); ASSERT(captured.size() == msg.size(), "hex blob frame size mismatch"); ASSERT(std::memcmp(captured.data(), msg.data(), msg.size()) == 0, "hex blob content mismatch"); return 1; } static int test_fe_move_semantics() { antpeer::FrameExtractor fe1; size_t count = 0; fe1.on_frame([&](const uint8_t*, size_t) { count++; }); auto frame = antheos::bus::ping("ABCD"); ASSERT(frame.has_value(), "frame build failed"); /* Move-construct */ antpeer::FrameExtractor fe2(std::move(fe1)); fe2.feed(frame->data(), frame->size()); ASSERT(count == 1, "moved extractor did not emit frame"); /* Move-assign */ antpeer::FrameExtractor fe3; fe3 = std::move(fe2); fe3.feed(frame->data(), frame->size()); ASSERT(count == 2, "move-assigned extractor did not emit frame"); return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Threaded Dispatch tests * ═══════════════════════════════════════════════════════════════════════ */ static int test_threaded_dispatch_basic() { auto sbus = test_membus_open("/antpeer_test_td1", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_td1", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; std::atomic calls{0}; s.on_request([&calls](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { calls++; std::string resp = "ECHO:" + std::string(cmd); srv.reply(sid, resp); }); s.set_dispatch_threads(2); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TdSvc1"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("TdSvc1", 3000); ASSERT(rc == 0, "discover failed"); char reply[256] = {}; int n = c.call("HELLO", 5000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); ASSERT(std::strcmp(reply, "ECHO:HELLO") == 0, "reply mismatch"); ASSERT(calls.load() == 1, "handler not called"); s.stop(); t.join(); return 1; } static int test_threaded_dispatch_concurrent() { auto sbus = test_membus_open("/antpeer_test_td2", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus != nullptr, "bus open failed"); antpeer::Server s; std::atomic calls{0}; s.on_request([&calls](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { calls++; std::string resp = "OK:" + std::string(cmd); srv.reply(sid, resp); }); s.set_dispatch_threads(2); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TdSvc2"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); /* 3 sequential calls from separate clients */ for (int i = 0; i < 3; i++) { auto cbus = test_membus_open("/antpeer_test_td2", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(cbus != nullptr, "client bus open failed"); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("TdSvc2", 3000); ASSERT(rc == 0, "discover failed"); std::string cmd = "REQ" + std::to_string(i); char reply[256] = {}; int n = c.call(cmd, 5000, reply, sizeof(reply)); ASSERT(n > 0, "call returned no data"); std::string expected = "OK:" + cmd; ASSERT(expected == reply, "reply mismatch"); } ASSERT(calls.load() == 3, "expected 3 handler calls"); s.stop(); t.join(); return 1; } static int test_threaded_dispatch_slow_handler() { auto sbus = test_membus_open("/antpeer_test_td3", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_td3", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; std::atomic calls{0}; s.on_request([&calls](antpeer::Server& srv, std::string_view sid, std::string_view cmd) { if (cmd == "SLOW") { usleep(200000); /* 200ms */ } calls++; srv.reply(sid, "DONE"); }); s.set_dispatch_threads(2); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TdSvc3"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "cli-oid", "cli-did", TEST_CLI_IID); ASSERT(rc == 0, "client init failed"); rc = c.discover("TdSvc3", 3000); ASSERT(rc == 0, "discover failed"); /* Send a slow request, then verify server is still responsive */ char reply[256] = {}; int n = c.call("SLOW", 5000, reply, sizeof(reply)); ASSERT(n > 0, "slow call returned no data"); ASSERT(std::strcmp(reply, "DONE") == 0, "slow reply mismatch"); ASSERT(calls.load() >= 1, "handler not called"); s.stop(); t.join(); return 1; } static int test_threaded_dispatch_stop() { auto sbus = test_membus_open("/antpeer_test_td4", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus != nullptr, "bus open failed"); antpeer::Server s; s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { usleep(50000); /* 50ms — simulate work */ srv.reply(sid, "OK"); }); s.set_dispatch_threads(2); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "TdSvc4"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); /* Let workers spin up, then stop immediately */ usleep(10000); s.stop(); t.join(); /* If we get here without deadlock or crash, the test passes */ return 1; } /* ═══════════════════════════════════════════════════════════════════════ * Main * ═══════════════════════════════════════════════════════════════════════ */ /* ═══════════════════════════════════════════════════════════════ * Auth (Z-verb) integration tests * ═══════════════════════════════════════════════════════════════ */ static int test_auth_success() { /* Generate keys: server trusts client's public key */ auto client_kp = antpeer::crypto::keygen(); /* Write client private key file */ char key_path[] = "/tmp/antpeer_test_key_XXXXXX"; int fd = mkstemp(key_path); close(fd); antpeer::crypto::save_key(client_kp, key_path); /* Write trusted keys file */ char trust_path[] = "/tmp/antpeer_test_trust_XXXXXX"; fd = mkstemp(trust_path); FILE* f = fdopen(fd, "w"); fprintf(f, "AUTH-ORG %s %s\n", client_kp.key_id.c_str(), client_kp.public_key.c_str()); fclose(f); auto sbus = test_membus_open("/antpeer_test_auth1", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_auth1", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; bool auth_ok = false; std::string auth_oid; s.on_auth([&](const antpeer::PeerIdentity& id, bool authenticated) { auth_ok = authenticated; auth_oid = id.oid; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { usleep(300000); /* give V+Z exchange time on poll thread */ srv.reply(sid, "OK"); }); int rc = s.require_auth(trust_path); ASSERT(rc == 0, "require_auth failed"); s.set_dispatch_threads(1); /* threaded: poll loop runs V+Z while handler waits */ rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "AuthSvc1"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "AUTH-ORG", "AUTH-DEV", "SN001"); ASSERT(rc == 0, "client init failed"); rc = c.set_auth_key(key_path); ASSERT(rc == 0, "set_auth_key failed"); rc = c.discover("AuthSvc1", 3000); ASSERT(rc == 0, "discover failed"); /* * call() pumps the bus — V+Z exchange happens during the pump. * Threaded dispatch ensures the server poll loop processes V+Z * while the request handler runs on a worker thread. * The handler delays briefly so auth completes before reply. */ char reply[64] = {}; c.call("PING", 5000, reply, sizeof(reply)); /* Wait for auth callback (async, fires after Z response processed) */ for (int i = 0; i < 50 && !auth_ok; ++i) usleep(50000); s.stop(); t.join(); unlink(key_path); unlink(trust_path); ASSERT(auth_ok, "auth should succeed"); ASSERT(auth_oid == "AUTH-ORG", "auth OID mismatch"); return 1; } static int test_auth_failure_wrong_key() { /* Client has key, but server trusts a DIFFERENT key */ auto client_kp = antpeer::crypto::keygen(); auto other_kp = antpeer::crypto::keygen(); char key_path[] = "/tmp/antpeer_test_key2_XXXXXX"; int fd = mkstemp(key_path); close(fd); antpeer::crypto::save_key(client_kp, key_path); char trust_path[] = "/tmp/antpeer_test_trust2_XXXXXX"; fd = mkstemp(trust_path); FILE* f = fdopen(fd, "w"); /* Trust the OTHER key, not the client's */ fprintf(f, "BAD-ORG %s %s\n", other_kp.key_id.c_str(), other_kp.public_key.c_str()); fclose(f); auto sbus = test_membus_open("/antpeer_test_auth2", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_auth2", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; bool auth_fired = false; bool auth_result = true; /* should become false */ s.on_auth([&](const antpeer::PeerIdentity&, bool authenticated) { auth_fired = true; auth_result = authenticated; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { usleep(300000); srv.reply(sid, "OK"); }); s.require_auth(trust_path); s.set_dispatch_threads(1); int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "AuthSvc2"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "BAD-ORG", "BAD-DEV", "SN001"); ASSERT(rc == 0, "client init failed"); c.set_auth_key(key_path); rc = c.discover("AuthSvc2", 3000); ASSERT(rc == 0, "discover failed"); char reply[64] = {}; c.call("PING", 5000, reply, sizeof(reply)); usleep(500000); s.stop(); t.join(); unlink(key_path); unlink(trust_path); /* The client's key_id doesn't match any trusted key, so the server can't find a matching pending_auth entry. auth_fn never fires. This is correct — unknown key_id = silent rejection. */ ASSERT(!auth_fired || !auth_result, "auth should not succeed with wrong key"); return 1; } static int test_auth_no_auth_backward_compat() { /* Server does NOT require auth. Client has no key. Should work fine. */ auto sbus = test_membus_open("/antpeer_test_auth3", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_auth3", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; bool verify_fired = false; s.on_verify([&](const antpeer::PeerIdentity&) { verify_fired = true; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { srv.reply(sid, "OK"); }); /* No require_auth() call */ int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "AuthSvc3"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "NOAUTH-ORG", "NOAUTH-DEV", "SN001"); ASSERT(rc == 0, "client init failed"); /* No set_auth_key() */ rc = c.discover("AuthSvc3", 3000); ASSERT(rc == 0, "discover failed"); char reply[64] = {}; int n = c.call("PING", 3000, reply, sizeof(reply)); s.stop(); t.join(); ASSERT(n > 0, "call should succeed without auth"); ASSERT(verify_fired, "verify should still fire"); return 1; } /* ═══════════════════════════════════════════════════════════════ * Relay Auth (multi-hop Z-verb) tests * ═══════════════════════════════════════════════════════════════ */ static int test_relay_auth_route_api() { /* Verify add_route/set_route store and don't crash */ auto sbus = test_membus_open("/antpeer_test_ra1", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus, "bus open failed"); antpeer::Server s; int rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "RouteSvc"); ASSERT(rc == 0, "server init failed"); /* add_route should not crash on valid args */ s.add_route("REMOTE", "AA.BB.CC", 0); antpeer::Client c; auto cbus = test_membus_open("/antpeer_test_ra1", antpeer::MEMBUS_DEFAULT_SIZE); rc = c.init(*cbus, "CLI-ORG", "CLI-DEV", "CL001"); ASSERT(rc == 0, "client init failed"); /* set_route should not crash */ c.set_route("AA.BB.CC", 2); s.stop(); return 1; } static int test_relay_auth_success() { /* Server knows client via a relay route. * When auth is required, server sends relay_auth_challenge. * Client has set_route, so it responds via relay_auth_response. * Same bus — the relay frames still parse and auth succeeds. */ auto client_kp = antpeer::crypto::keygen(); char key_path[] = "/tmp/antpeer_test_rkey_XXXXXX"; int fd = mkstemp(key_path); close(fd); antpeer::crypto::save_key(client_kp, key_path); char trust_path[] = "/tmp/antpeer_test_rtrust_XXXXXX"; fd = mkstemp(trust_path); FILE* f = fdopen(fd, "w"); fprintf(f, "RELAY-ORG %s %s\n", client_kp.key_id.c_str(), client_kp.public_key.c_str()); fclose(f); auto sbus = test_membus_open("/antpeer_test_ra2", antpeer::MEMBUS_DEFAULT_SIZE); auto cbus = test_membus_open("/antpeer_test_ra2", antpeer::MEMBUS_DEFAULT_SIZE); ASSERT(sbus && cbus, "bus open failed"); antpeer::Server s; bool auth_ok = false; std::string auth_oid; s.on_auth([&](const antpeer::PeerIdentity& id, bool authenticated) { auth_ok = authenticated; auth_oid = id.oid; }); s.on_request([](antpeer::Server& srv, std::string_view sid, std::string_view) { usleep(300000); srv.reply(sid, "OK"); }); int rc = s.require_auth(trust_path); ASSERT(rc == 0, "require_auth failed"); s.set_dispatch_threads(1); rc = s.init(*sbus, TEST_OID, TEST_DID, TEST_IID, "RelaySvc"); ASSERT(rc == 0, "server init failed"); std::thread t([&s]() { s.run(); }); antpeer::Client c; rc = c.init(*cbus, "RELAY-ORG", "RELAY-DEV", "SN001"); ASSERT(rc == 0, "client init failed"); rc = c.set_auth_key(key_path); ASSERT(rc == 0, "set_auth_key failed"); rc = c.discover("RelaySvc", 3000); ASSERT(rc == 0, "discover failed"); /* Now add relay routes. The server's BID and client's BID are known * after discover. Build a simple 2-hop path: SERVER.CLIENT */ const char* srv_bid = s.bid(); const char* cli_bid = c.peer_bid(); ASSERT(srv_bid && cli_bid, "BIDs not available"); /* Path: server_bid.client_bid — server at index 0, client at index 1 */ std::string path = std::string(srv_bid) + "." + std::string(cli_bid); s.add_route(cli_bid, path, 0); c.set_route(path, 1); char reply[64] = {}; c.call("PING", 5000, reply, sizeof(reply)); for (int i = 0; i < 50 && !auth_ok; ++i) usleep(50000); s.stop(); t.join(); unlink(key_path); unlink(trust_path); ASSERT(auth_ok, "relay auth should succeed"); ASSERT(auth_oid == "RELAY-ORG", "relay auth OID mismatch"); return 1; } int main() { ShmSweeper sweeper; /* unlinks every bus this run opened, at exit */ std::fprintf(stderr, "\n=== AntPeer Unit Tests (C++17) ===\n\n"); std::fprintf(stderr, "--- Bus ---\n"); RUN_TEST(test_bus_membus_open_close); RUN_TEST(test_bus_write_read); RUN_TEST(test_bus_reopen); RUN_TEST(test_bus_read_empty); RUN_TEST(test_bus_empty_name); RUN_TEST(test_bus_zero_size); RUN_TEST(test_bus_factory_throws); RUN_TEST(test_bus_large_write_read); RUN_TEST(test_bus_multiple_writes); RUN_TEST(test_bus_reopen_preserves_name); std::fprintf(stderr, "\n--- Server ---\n"); RUN_TEST(test_server_create_destroy); RUN_TEST(test_server_init); RUN_TEST(test_server_establish); RUN_TEST(test_server_stop); RUN_TEST(test_server_broadcast); RUN_TEST(test_server_handler_reply); RUN_TEST(test_server_session_handler); RUN_TEST(test_server_default_state); RUN_TEST(test_server_init_bad_args); RUN_TEST(test_server_double_destroy); RUN_TEST(test_server_reply_unknown_sid); RUN_TEST(test_server_finish_unknown_sid); RUN_TEST(test_server_offer_additional); RUN_TEST(test_server_bid_before_init); RUN_TEST(test_server_set_expiry); RUN_TEST(test_server_notify_unknown_sid); RUN_TEST(test_server_multiple_broadcasts); RUN_TEST(test_server_broadcast_empty_event); std::fprintf(stderr, "\n--- Client ---\n"); RUN_TEST(test_client_create_destroy); RUN_TEST(test_client_init); RUN_TEST(test_client_discover); RUN_TEST(test_client_call_timeout); RUN_TEST(test_client_fire); RUN_TEST(test_client_close_reopen); RUN_TEST(test_client_notify); RUN_TEST(test_client_multi_call); RUN_TEST(test_client_required_shutdown); RUN_TEST(test_client_default_state); RUN_TEST(test_client_init_bad_args); RUN_TEST(test_client_double_destroy); RUN_TEST(test_client_call_not_connected); RUN_TEST(test_client_fire_not_connected); RUN_TEST(test_client_peer_bid_before_discover); RUN_TEST(test_client_discover_timeout); RUN_TEST(test_client_discover_empty_service); RUN_TEST(test_client_state_transitions); RUN_TEST(test_client_call_null_command); RUN_TEST(test_client_fire_null_command); std::fprintf(stderr, "\n--- Integration ---\n"); RUN_TEST(test_server_finish_handler_called); RUN_TEST(test_multi_client_same_server); RUN_TEST(test_call_large_payload); RUN_TEST(test_call_empty_command); RUN_TEST(test_client_call_small_reply_buffer); RUN_TEST(test_client_call_null_reply_buffer); std::fprintf(stderr, "\n--- FrameExtractor ---\n"); RUN_TEST(test_fe_create_destroy); RUN_TEST(test_fe_simple_frame); RUN_TEST(test_fe_frame_with_tail); RUN_TEST(test_fe_tail_with_som_eom_bytes); RUN_TEST(test_fe_multiple_frames); RUN_TEST(test_fe_partial_feed); RUN_TEST(test_fe_partial_feed_with_tail); RUN_TEST(test_fe_garbage_before_som); RUN_TEST(test_fe_large_tail); RUN_TEST(test_fe_multiple_blob_words); RUN_TEST(test_fe_reset); RUN_TEST(test_fe_no_callback); RUN_TEST(test_fe_tail_then_next_frame); RUN_TEST(test_fe_hex_blob); RUN_TEST(test_fe_move_semantics); std::fprintf(stderr, "\n--- Misc ---\n"); RUN_TEST(test_iid_caller_provided); std::fprintf(stderr, "\n--- Identity ---\n"); RUN_TEST(test_peer_identity_struct); RUN_TEST(test_accept_callback); RUN_TEST(test_verify_roundtrip); RUN_TEST(test_server_responds_to_verify); RUN_TEST(test_unverified_peer_no_identity); RUN_TEST(test_sid_peer_resolution); RUN_TEST(test_multiple_peers_verify); RUN_TEST(test_send_verify_explicit); std::fprintf(stderr, "\n--- Threaded Dispatch ---\n"); RUN_TEST(test_threaded_dispatch_basic); RUN_TEST(test_threaded_dispatch_concurrent); RUN_TEST(test_threaded_dispatch_slow_handler); RUN_TEST(test_threaded_dispatch_stop); std::fprintf(stderr, "\n--- Auth (Z-verb) ---\n"); RUN_TEST(test_auth_success); RUN_TEST(test_auth_failure_wrong_key); RUN_TEST(test_auth_no_auth_backward_compat); std::fprintf(stderr, "\n--- Relay Auth (multi-hop Z-verb) ---\n"); RUN_TEST(test_relay_auth_route_api); RUN_TEST(test_relay_auth_success); std::fprintf(stderr, "\n=== Results: %d/%d passed", tests_passed, tests_run); if (tests_failed > 0) std::fprintf(stderr, " (%d FAILED)", tests_failed); std::fprintf(stderr, " ===\n\n"); return tests_failed > 0 ? 1 : 0; }