@@ -505,4 +505,218 @@ TEST(ShmTest, MpscReactorPipelinedConcurrencyAndOrder)
505505 server->close ();
506506}
507507
508+ TEST (ShmTest, MpscReactorMultiClientResponseRouting)
509+ {
510+ constexpr uint32_t NUM_CLIENTS = 4 ;
511+ constexpr uint32_t K = 128 ;
512+ constexpr size_t RING_SIZE = 64UL * 1024 ;
513+
514+ std::string base_name = " shm_mpsc_multi_" + std::to_string (getpid ());
515+ auto server = IpcServer::create_mpsc_shm (base_name, NUM_CLIENTS , RING_SIZE , RING_SIZE );
516+ ASSERT_TRUE (server->listen ()) << " MPSC multi-client server failed to listen" ;
517+
518+ ReactorTestPool pool (8 );
519+ std::thread server_thread ([&]() {
520+ server->run_reactor ([&pool](int , std::span<const uint8_t > req, IpcServer::Respond respond) {
521+ std::vector<uint8_t > r (req.begin (), req.end ());
522+ pool.enqueue ([r = std::move (r), respond = std::move (respond)]() mutable {
523+ uint32_t seq = 0 ;
524+ std::memcpy (&seq, r.data () + sizeof (uint32_t ), sizeof (uint32_t ));
525+ std::this_thread::sleep_for (std::chrono::microseconds (20 * (seq % 8 )));
526+ respond (std::move (r));
527+ });
528+ });
529+ });
530+ std::this_thread::sleep_for (std::chrono::milliseconds (100 ));
531+
532+ std::atomic<int > wrong_client{ 0 };
533+ std::atomic<int > out_of_order{ 0 };
534+ std::atomic<int > stalls{ 0 };
535+
536+ auto run_client = [&](uint32_t c) {
537+ auto client = IpcClient::create_mpsc_shm (base_name, c);
538+ if (!client->connect ()) {
539+ stalls++;
540+ return ;
541+ }
542+ for (uint32_t s = 0 ; s < K; s++) {
543+ uint32_t msg[2 ] = { c, s };
544+ while (!client->send (msg, sizeof (msg), 100'000'000ULL )) {
545+ }
546+ }
547+ for (uint32_t s = 0 ; s < K; s++) {
548+ std::span<const uint8_t > resp;
549+ size_t empties = 0 ;
550+ while ((resp = client->receive (100'000'000ULL )).empty ()) {
551+ if (++empties > 50 ) {
552+ stalls++;
553+ break ;
554+ }
555+ }
556+ if (resp.empty ()) {
557+ break ;
558+ }
559+ uint32_t got[2 ] = { 0 , 0 };
560+ std::memcpy (got, resp.data (), sizeof (got));
561+ if (got[0 ] != c) {
562+ wrong_client++;
563+ }
564+ if (got[1 ] != s) {
565+ out_of_order++;
566+ }
567+ client->release (resp.size ());
568+ }
569+ client->close ();
570+ };
571+
572+ std::vector<std::thread> clients;
573+ for (uint32_t c = 0 ; c < NUM_CLIENTS ; c++) {
574+ clients.emplace_back (run_client, c);
575+ }
576+ for (auto & t : clients) {
577+ t.join ();
578+ }
579+
580+ server->request_shutdown ();
581+ server_thread.join ();
582+ server->close ();
583+
584+ EXPECT_EQ (wrong_client.load (), 0 ) << " responses were routed to the wrong client" ;
585+ EXPECT_EQ (out_of_order.load (), 0 ) << " responses arrived out of order within a client" ;
586+ EXPECT_EQ (stalls.load (), 0 ) << " a client stalled waiting for its responses" ;
587+ }
588+
589+ TEST (ShmTest, MpscReactorAutoClaimMultiClient)
590+ {
591+ constexpr uint32_t NUM_CLIENTS = 4 ;
592+ constexpr uint32_t K = 128 ;
593+ constexpr size_t RING_SIZE = 64UL * 1024 ;
594+
595+ std::string base_name = " shm_mpsc_autoclaim_" + std::to_string (getpid ());
596+ auto server = IpcServer::create_mpsc_shm (base_name, NUM_CLIENTS , RING_SIZE , RING_SIZE );
597+ ASSERT_TRUE (server->listen ()) << " MPSC auto-claim server failed to listen" ;
598+
599+ ReactorTestPool pool (8 );
600+ std::thread server_thread ([&]() {
601+ server->run_reactor ([&pool](int , std::span<const uint8_t > req, IpcServer::Respond respond) {
602+ std::vector<uint8_t > r (req.begin (), req.end ());
603+ pool.enqueue ([r = std::move (r), respond = std::move (respond)]() mutable {
604+ uint32_t seq = 0 ;
605+ std::memcpy (&seq, r.data () + sizeof (uint32_t ), sizeof (uint32_t ));
606+ std::this_thread::sleep_for (std::chrono::microseconds (20 * (seq % 8 )));
607+ respond (std::move (r));
608+ });
609+ });
610+ });
611+ std::this_thread::sleep_for (std::chrono::milliseconds (100 ));
612+
613+ std::atomic<int > wrong_client{ 0 };
614+ std::atomic<int > out_of_order{ 0 };
615+ std::atomic<int > stalls{ 0 };
616+
617+ auto run_client = [&](uint32_t c) {
618+ auto client = IpcClient::create_mpsc_shm (base_name);
619+ if (!client->connect ()) {
620+ stalls++;
621+ return ;
622+ }
623+ for (uint32_t s = 0 ; s < K; s++) {
624+ uint32_t msg[2 ] = { c, s };
625+ while (!client->send (msg, sizeof (msg), 100'000'000ULL )) {
626+ }
627+ }
628+ for (uint32_t s = 0 ; s < K; s++) {
629+ std::span<const uint8_t > resp;
630+ size_t empties = 0 ;
631+ while ((resp = client->receive (100'000'000ULL )).empty ()) {
632+ if (++empties > 50 ) {
633+ stalls++;
634+ break ;
635+ }
636+ }
637+ if (resp.empty ()) {
638+ break ;
639+ }
640+ uint32_t got[2 ] = { 0 , 0 };
641+ std::memcpy (got, resp.data (), sizeof (got));
642+ if (got[0 ] != c) {
643+ wrong_client++;
644+ }
645+ if (got[1 ] != s) {
646+ out_of_order++;
647+ }
648+ client->release (resp.size ());
649+ }
650+ client->close ();
651+ };
652+
653+ std::vector<std::thread> clients;
654+ for (uint32_t c = 0 ; c < NUM_CLIENTS ; c++) {
655+ clients.emplace_back (run_client, c);
656+ }
657+ for (auto & t : clients) {
658+ t.join ();
659+ }
660+
661+ server->request_shutdown ();
662+ server_thread.join ();
663+ server->close ();
664+
665+ EXPECT_EQ (wrong_client.load (), 0 ) << " self-allocated clients aliased onto a shared slot" ;
666+ EXPECT_EQ (out_of_order.load (), 0 ) << " responses arrived out of order within a client" ;
667+ EXPECT_EQ (stalls.load (), 0 ) << " a client stalled" ;
668+ }
669+
670+ TEST (ShmTest, MpscSlotFreedOnDisconnectAndReclaimed)
671+ {
672+ constexpr size_t NUM_CLIENTS = 1 ;
673+ constexpr size_t RING_SIZE = 16UL * 1024 ;
674+
675+ std::string base_name = " shm_mpsc_reuse_" + std::to_string (getpid ());
676+ auto server = IpcServer::create_mpsc_shm (base_name, NUM_CLIENTS , RING_SIZE , RING_SIZE );
677+ ASSERT_TRUE (server->listen ()) << " server failed to listen" ;
678+
679+ ReactorTestPool pool (2 );
680+ std::thread server_thread ([&]() {
681+ server->run_reactor ([&pool](int , std::span<const uint8_t > req, IpcServer::Respond respond) {
682+ std::vector<uint8_t > r (req.begin (), req.end ());
683+ pool.enqueue ([r = std::move (r), respond = std::move (respond)]() mutable { respond (std::move (r)); });
684+ });
685+ });
686+ std::this_thread::sleep_for (std::chrono::milliseconds (100 ));
687+
688+ auto round_trip = [](IpcClient& client, uint32_t value) -> bool {
689+ if (!client.send (&value, sizeof (value), 100'000'000ULL )) {
690+ return false ;
691+ }
692+ std::span<const uint8_t > resp;
693+ size_t empties = 0 ;
694+ while ((resp = client.receive (100'000'000ULL )).empty ()) {
695+ if (++empties > 50 ) {
696+ return false ;
697+ }
698+ }
699+ uint32_t got = 0 ;
700+ std::memcpy (&got, resp.data (), sizeof (got));
701+ client.release (resp.size ());
702+ return got == value;
703+ };
704+
705+ {
706+ auto a = IpcClient::create_mpsc_shm (base_name);
707+ ASSERT_TRUE (a->connect ()) << " client A failed to connect" ;
708+ EXPECT_TRUE (round_trip (*a, 0x1111 )) << " client A round-trip failed" ;
709+ a->close ();
710+ }
711+
712+ auto b = IpcClient::create_mpsc_shm (base_name);
713+ ASSERT_TRUE (b->connect ()) << " client B failed to reclaim the freed slot" ;
714+ EXPECT_TRUE (round_trip (*b, 0x2222 )) << " client B round-trip over the reused ring failed" ;
715+ b->close ();
716+
717+ server->request_shutdown ();
718+ server_thread.join ();
719+ server->close ();
720+ }
721+
508722} // namespace
0 commit comments