From 6dbebcacca27d6c8f57fe23eb69826edb9e2a060 Mon Sep 17 00:00:00 2001 From: Venkit Kasiviswanathan Date: Tue, 12 May 2026 00:49:01 +0000 Subject: [PATCH] [orchagent]: Add ZmqRouteServer for concurrent route updates Introduce a dedicated ZmqRouteServer (and ZmqRouteOrch/ZmqRouteConsumer) used by RouteOrch for receiving APPL_DB route updates from fpmsyncd. Non-fabric/non-dpu switches now create a ZmqRouteServer instead of the generic ZmqServer; fabric and DPU continue to use ZmqServer. ZmqRouteConsumer merges incoming tuples into m_toSync from the mqPollThread ingress callback under m_toSyncMutex, and only notifies the orch main loop once a batch (gMaxBulkSize) has accumulated. To support this concurrent access, ConsumerBase::addToSync and dumpPendingTasks are made virtual so the route consumer can wrap them in a lock, while the default single-threaded base remains lock-free. Also: - create_zmq_route_server() factory added in lib/orch_zmq_config. - getCfgSwitchType() moved earlier in orchagent main so the server type can be chosen based on switch type. - Update fake_zmqserver to return a ZmqMessageHandler* from handleReceivedData to match the new upstream signature. - Add zmq_route_orch_ut.cpp unit tests. Signed-off-by: Venkit Kasiviswanathan --- lib/orch_zmq_config.cpp | 16 + lib/orch_zmq_config.h | 2 + orchagent/Makefile.am | 1 + orchagent/main.cpp | 11 +- orchagent/orch.cpp | 4 + orchagent/orch.h | 14 +- orchagent/orchdaemon.cpp | 2 +- orchagent/p4orch/tests/Makefile.am | 1 + orchagent/p4orch/tests/fake_zmqserver.cpp | 2 +- orchagent/routeorch.cpp | 4 +- orchagent/routeorch.h | 8 +- orchagent/zmqrouteorch.cpp | 114 +++++++ orchagent/zmqrouteorch.h | 49 +++ tests/mock_tests/Makefile.am | 2 + tests/mock_tests/zmq_route_orch_ut.cpp | 359 ++++++++++++++++++++++ 15 files changed, 573 insertions(+), 16 deletions(-) create mode 100644 orchagent/zmqrouteorch.cpp create mode 100644 orchagent/zmqrouteorch.h create mode 100644 tests/mock_tests/zmq_route_orch_ut.cpp diff --git a/lib/orch_zmq_config.cpp b/lib/orch_zmq_config.cpp index 09bc66e0b0c..2775e1981c4 100644 --- a/lib/orch_zmq_config.cpp +++ b/lib/orch_zmq_config.cpp @@ -78,6 +78,22 @@ std::shared_ptr swss::create_zmq_server(std::string zmq_address return std::make_shared(zmq_address, vrf, true); } +std::shared_ptr swss::create_zmq_route_server(std::string zmq_address, std::string vrf) +{ + if (!std::regex_search(zmq_address, ZMQ_NONE_IPV6_ADDRESS_WITH_PORT) + && !std::regex_search(zmq_address, ZMQ_IPV6_ADDRESS_WITH_PORT)) + { + auto zmq_port = get_zmq_port(); + zmq_address = zmq_address + ":" + std::to_string(zmq_port); + } + + SWSS_LOG_NOTICE("Create ZMQ server with address: %s", zmq_address.c_str()); + + // To prevent message loss between ZmqServer's bind operation and the creation of ZmqProducerStateTable, + // use lazy binding and call bind() only after the handler has been registered. + return std::make_shared(zmq_address, vrf, true); +} + bool swss::get_feature_status(std::string feature, bool default_value) { std::shared_ptr enabled = nullptr; diff --git a/lib/orch_zmq_config.h b/lib/orch_zmq_config.h index 68aff440dbd..93da34cdcfb 100644 --- a/lib/orch_zmq_config.h +++ b/lib/orch_zmq_config.h @@ -9,6 +9,7 @@ #include "zmqclient.h" #include "zmqserver.h" #include "zmqproducerstatetable.h" +#include "zmqrouteserver.h" /* * swssconfig will only connect to local orchagent ZMQ endpoint. @@ -34,6 +35,7 @@ int get_zmq_port(); std::shared_ptr create_zmq_client(std::string zmq_address, std::string vrf=""); std::shared_ptr create_zmq_server(std::string zmq_address, std::string vrf=""); +std::shared_ptr create_zmq_route_server(std::string zmq_address, std::string vrf=""); bool get_feature_status(std::string feature, bool default_value); diff --git a/orchagent/Makefile.am b/orchagent/Makefile.am index 9a372674751..ae0f5509705 100644 --- a/orchagent/Makefile.am +++ b/orchagent/Makefile.am @@ -117,6 +117,7 @@ orchagent_SOURCES = \ response_publisher.cpp \ nvgreorch.cpp \ zmqorch.cpp \ + zmqrouteorch.cpp \ dash/dashenifwdorch.cpp \ dash/dashenifwdinfo.cpp \ dash/dashcounter.cpp \ diff --git a/orchagent/main.cpp b/orchagent/main.cpp index 702064a5f37..cc864ba5b56 100644 --- a/orchagent/main.cpp +++ b/orchagent/main.cpp @@ -642,6 +642,9 @@ int main(int argc, char **argv) DBConnector config_db("CONFIG_DB", 0); DBConnector state_db("STATE_DB", 0); + // Get switch_type + getCfgSwitchType(&config_db, gMySwitchType, gMySwitchSubType); + // Instantiate ZMQ server shared_ptr zmq_server = nullptr; if (zmq_server_address.empty()) @@ -651,12 +654,12 @@ int main(int argc, char **argv) else { SWSS_LOG_NOTICE("The ZMQ channel on the northbound side of orchagent has been initialized: %s, %s", zmq_server_address.c_str(), vrf.c_str()); - zmq_server = create_zmq_server(zmq_server_address); + if (gMySwitchType == "fabric" || gMySwitchType == "dpu") + zmq_server = create_zmq_server(zmq_server_address); + else + zmq_server = create_zmq_route_server(zmq_server_address); } - // Get switch_type - getCfgSwitchType(&config_db, gMySwitchType, gMySwitchSubType); - sai_attribute_t attr; vector attrs; diff --git a/orchagent/orch.cpp b/orchagent/orch.cpp index 2a9e72b03a0..1ed99251264 100644 --- a/orchagent/orch.cpp +++ b/orchagent/orch.cpp @@ -411,6 +411,10 @@ size_t ConsumerBase::addToSync(const std::deque &entries recordTuples(entries); } + // Call addToSyncInternal directly so we don't re-enter virtual dispatch + // (and any subclass-installed lock) per entry. Subclasses that need + // locking override addToSync(deque) to take the lock once before calling + // this base implementation. for (auto& entry: entries) { addToSyncInternal(entry, onRetry, onRetry); diff --git a/orchagent/orch.h b/orchagent/orch.h index 35b79ef65fc..8de5ad2452f 100644 --- a/orchagent/orch.h +++ b/orchagent/orch.h @@ -167,7 +167,13 @@ class ConsumerBase : public Executor { } std::string dumpTuple(const swss::KeyOpFieldsValuesTuple &tuple); - void dumpPendingTasks(std::vector &ts); + + /* + * dumpPendingTasks and the addToSync overloads are virtual so concurrent + * subclasses (e.g. ZmqRouteConsumer) can wrap the base implementation in + * a lock. The base class itself is single-threaded and takes no lock. + */ + virtual void dumpPendingTasks(std::vector &ts); /* Store the latest 'golden' status */ // TODO: hide? @@ -177,11 +183,11 @@ class ConsumerBase : public Executor { void recordTuple(const swss::KeyOpFieldsValuesTuple &tuple); void recordTuples(const std::deque &entries); - void addToSync(const swss::KeyOpFieldsValuesTuple &entry, bool onRetry=false); + virtual void addToSync(const swss::KeyOpFieldsValuesTuple &entry, bool onRetry=false); // Returns: the number of entries added to m_toSync - size_t addToSync(const std::deque &entries, bool onRetry=false); - size_t addToSync(std::shared_ptr> entries, bool onRetry=false); + virtual size_t addToSync(const std::deque &entries, bool onRetry=false); + size_t addToSync(std::shared_ptr> entries, bool onRetry=false); /** * @brief Add the failed task and its constraint to the consumer's RetryCache diff --git a/orchagent/orchdaemon.cpp b/orchagent/orchdaemon.cpp index 4c14dad7097..e7623f2f58a 100644 --- a/orchagent/orchdaemon.cpp +++ b/orchagent/orchdaemon.cpp @@ -332,7 +332,7 @@ bool OrchDaemon::init() // Enable the fpmsyncd service to send Route events to orchagent via the ZMQ channel. auto enable_route_zmq = get_feature_status(ORCH_NORTHBOND_ROUTE_ZMQ_ENABLED, false); - auto route_zmq_sever = enable_route_zmq ? m_zmqServer : nullptr; + auto route_zmq_sever = enable_route_zmq ? dynamic_cast(m_zmqServer) : nullptr; gRouteOrch = new RouteOrch(m_applDb, route_tables, gSwitchOrch, gNeighOrch, gIntfsOrch, vrf_orch, gFgNhgOrch, gSrv6Orch, route_zmq_sever); gNhgOrch = new NhgOrch(m_applDb, APP_NEXTHOP_GROUP_TABLE_NAME); diff --git a/orchagent/p4orch/tests/Makefile.am b/orchagent/p4orch/tests/Makefile.am index cfa3881e264..30e0676dfd6 100644 --- a/orchagent/p4orch/tests/Makefile.am +++ b/orchagent/p4orch/tests/Makefile.am @@ -29,6 +29,7 @@ p4orch_tests_SOURCES = $(ORCHAGENT_DIR)/orch.cpp \ $(ORCHAGENT_DIR)/request_parser.cpp \ $(top_srcdir)/lib/recorder.cpp \ $(ORCHAGENT_DIR)/zmqorch.cpp \ + $(ORCHAGENT_DIR)/zmqrouteorch.cpp \ $(ORCHAGENT_DIR)/flex_counter/flex_counter_manager.cpp \ $(ORCHAGENT_DIR)/flex_counter/flow_counter_handler.cpp \ $(ORCHAGENT_DIR)/port/port_capabilities.cpp \ diff --git a/orchagent/p4orch/tests/fake_zmqserver.cpp b/orchagent/p4orch/tests/fake_zmqserver.cpp index 505a5f0acee..64f8d47d8c6 100644 --- a/orchagent/p4orch/tests/fake_zmqserver.cpp +++ b/orchagent/p4orch/tests/fake_zmqserver.cpp @@ -22,7 +22,7 @@ ZmqMessageHandler* ZmqServer::findMessageHandler(const std::string dbName, return nullptr; } -void ZmqServer::handleReceivedData(const char* buffer, const size_t size) {} +ZmqMessageHandler* ZmqServer::handleReceivedData(const char* buffer, const size_t size) { return nullptr; } void ZmqServer::mqPollThread() {} diff --git a/orchagent/routeorch.cpp b/orchagent/routeorch.cpp index b8fba72469f..bde2a79bcc2 100644 --- a/orchagent/routeorch.cpp +++ b/orchagent/routeorch.cpp @@ -37,11 +37,11 @@ extern string gMySwitchType; #define DEFAULT_NUMBER_OF_ECMP_GROUPS 128 #define DEFAULT_MAX_ECMP_GROUP_SIZE 32 -RouteOrch::RouteOrch(DBConnector *db, vector &tableNames, SwitchOrch *switchOrch, NeighOrch *neighOrch, IntfsOrch *intfsOrch, VRFOrch *vrfOrch, FgNhgOrch *fgNhgOrch, Srv6Orch *srv6Orch, swss::ZmqServer *zmqServer) : +RouteOrch::RouteOrch(DBConnector *db, vector &tableNames, SwitchOrch *switchOrch, NeighOrch *neighOrch, IntfsOrch *intfsOrch, VRFOrch *vrfOrch, FgNhgOrch *fgNhgOrch, Srv6Orch *srv6Orch, ZmqRouteServer *zmqRouteServer) : gRouteBulker(sai_route_api, gMaxBulkSize), gLabelRouteBulker(sai_mpls_api, gMaxBulkSize), gNextHopGroupMemberBulker(sai_next_hop_group_api, gSwitchId, gMaxBulkSize), - ZmqOrch(db, tableNames, zmqServer), + ZmqRouteOrch(db, tableNames, zmqRouteServer), m_switchOrch(switchOrch), m_neighOrch(neighOrch), m_intfsOrch(intfsOrch), diff --git a/orchagent/routeorch.h b/orchagent/routeorch.h index 5fdb5b8e462..8f9b458d694 100644 --- a/orchagent/routeorch.h +++ b/orchagent/routeorch.h @@ -16,8 +16,8 @@ #include "bulker.h" #include "fgnhgorch.h" #include -#include "zmqorch.h" -#include "zmqserver.h" +#include "zmqrouteorch.h" +#include "zmqrouteserver.h" #include /* Maximum next hop group number */ @@ -212,10 +212,10 @@ struct LabelRouteBulkContext } }; -class RouteOrch : public ZmqOrch, public Subject +class RouteOrch : public ZmqRouteOrch, public Subject { public: - RouteOrch(DBConnector *db, vector &tableNames, SwitchOrch *switchOrch, NeighOrch *neighOrch, IntfsOrch *intfsOrch, VRFOrch *vrfOrch, FgNhgOrch *fgNhgOrch, Srv6Orch *srv6Orch, swss::ZmqServer *zmqServer = nullptr); + RouteOrch(DBConnector *db, vector &tableNames, SwitchOrch *switchOrch, NeighOrch *neighOrch, IntfsOrch *intfsOrch, VRFOrch *vrfOrch, FgNhgOrch *fgNhgOrch, Srv6Orch *srv6Orch, ZmqRouteServer *zmqServer = nullptr); bool hasNextHopGroup(const NextHopGroupKey&) const; sai_object_id_t getNextHopGroupId(const NextHopGroupKey&); diff --git a/orchagent/zmqrouteorch.cpp b/orchagent/zmqrouteorch.cpp new file mode 100644 index 00000000000..fdb2a4be9c8 --- /dev/null +++ b/orchagent/zmqrouteorch.cpp @@ -0,0 +1,114 @@ +#include "zmqrouteorch.h" + +using namespace swss; +using namespace std; + +extern int gBatchSize; +extern size_t gMaxBulkSize; + +ZmqRouteConsumer::ZmqRouteConsumer(ZmqRouteConsumerStateTable *select, Orch *orch, const std::string &name) + : ConsumerBase(select, orch, name) +{ + // mqPollThread runs the merge inline: kcos go straight into m_toSync + // under m_toSyncMutex. The eventfd is fired only when m_toSync grows past + // gMaxBulkSize (so the orch main loop has a real batch to drain); + // otherwise mqPollThread fires it once per burst after the burst quiesces. + select->setIngressCallback( + [this, select](const std::vector> &kcos) { + std::lock_guard lk(m_toSyncMutex); + for (const auto &kco : kcos) + { + // Qualified call to bypass our own virtual override (which + // would re-acquire m_toSyncMutex per entry). + ConsumerBase::addToSync(*kco, /*onRetry=*/false); + } + if (m_toSync.size() >= gMaxBulkSize) + { + select->notifyPending(); + } + }); +} + +void ZmqRouteConsumer::execute() +{ + SWSS_LOG_ENTER(); + + // Tuples were already merged into m_toSync by the ingress callback running + // on mqPollThread. The main loop's job is just to drain. + drain(); +} + +void ZmqRouteConsumer::drain() +{ + std::lock_guard lk(m_toSyncMutex); + if (!m_toSync.empty()) + (static_cast(m_orch))->doTask(*this); +} + +void ZmqRouteConsumer::addToSync(const KeyOpFieldsValuesTuple &entry, bool onRetry) +{ + std::lock_guard lk(m_toSyncMutex); + ConsumerBase::addToSync(entry, onRetry); +} + +size_t ZmqRouteConsumer::addToSync(const std::deque &entries, bool onRetry) +{ + std::lock_guard lk(m_toSyncMutex); + return ConsumerBase::addToSync(entries, onRetry); +} + +void ZmqRouteConsumer::dumpPendingTasks(std::vector &ts) +{ + std::lock_guard lk(m_toSyncMutex); + ConsumerBase::dumpPendingTasks(ts); +} + + +ZmqRouteOrch::ZmqRouteOrch(DBConnector *db, const vector &tableNames, ZmqRouteServer *zmqServer) +: Orch() +{ + for (auto it : tableNames) + { + addConsumer(db, it, default_orch_pri, zmqServer); + } +} + + +ZmqRouteOrch::ZmqRouteOrch(DBConnector *db, const vector &tableNames_with_pri, ZmqRouteServer *zmqServer) +{ + for (const auto& it : tableNames_with_pri) + { + addConsumer(db, it.first, it.second, zmqServer); + } +} + +void ZmqRouteOrch::addConsumer(DBConnector *db, string tableName, int pri, ZmqRouteServer *zmqServer) +{ + if (db->getDbId() == APPL_DB || db->getDbId() == DPU_APPL_DB) + { + if (zmqServer != nullptr) + { + SWSS_LOG_DEBUG("ZmqRouteConsumer initialize for: %s", tableName.c_str()); + addExecutor( + new ZmqRouteConsumer( + new ZmqRouteConsumerStateTable( + db, tableName, *zmqServer, pri, /* dbPersistence= */false), + this, tableName)); + } + else + { + SWSS_LOG_DEBUG("Consumer initialize for: %s", tableName.c_str()); + addExecutor(new Consumer(new ConsumerStateTable(db, tableName, gBatchSize, pri), this, tableName)); + } + } + else + { + SWSS_LOG_WARN("ZmqRouteOrch does not support create consumer for db: %d, table: %s", db->getDbId(), tableName.c_str()); + } +} + +void ZmqRouteOrch::doTask(Consumer &consumer) +{ + // When ZMQ disabled, forward data from Consumer + doTask((ConsumerBase &)consumer); +} diff --git a/orchagent/zmqrouteorch.h b/orchagent/zmqrouteorch.h new file mode 100644 index 00000000000..aa3169007c4 --- /dev/null +++ b/orchagent/zmqrouteorch.h @@ -0,0 +1,49 @@ +#pragma once + +#include +#include +#include +#include +#include +#include +#include "zmqrouteserver.h" +#include "zmqrouteconsumerstatetable.h" + +extern int gZmqExecuteTimeQuantaMsecs; + +class ZmqRouteConsumer : public ConsumerBase { +public: + ZmqRouteConsumer(ZmqRouteConsumerStateTable *select, Orch *orch, const std::string &name); + + swss::TableBase *getConsumerTable() const override + { + // ZmqRouteConsumerStateTable is a subclass of TableBase + return static_cast(getSelectable()); + } + + void execute() override; + void drain() override; + + // Locked overrides: take m_toSyncMutex, then forward to ConsumerBase. + // These guard against the ZmqRouteServer mqPollThread (calling addToSync + // via the ingress callback) racing with the orch main thread. + void addToSync(const swss::KeyOpFieldsValuesTuple &entry, bool onRetry=false) override; + size_t addToSync(const std::deque &entries, bool onRetry=false) override; + void dumpPendingTasks(std::vector &ts) override; + +private: + mutable std::mutex m_toSyncMutex; +}; + +class ZmqRouteOrch : public Orch +{ +public: + ZmqRouteOrch(swss::DBConnector *db, const std::vector &tableNames, ZmqRouteServer *zmqServer); + ZmqRouteOrch(swss::DBConnector *db, const std::vector &tableNames_with_pri, ZmqRouteServer *zmqServer); + + virtual void doTask(ConsumerBase &consumer) { }; + void doTask(Consumer &consumer) override; + +private: + void addConsumer(swss::DBConnector *db, std::string tableName, int pri, ZmqRouteServer *zmqServer); +}; diff --git a/tests/mock_tests/Makefile.am b/tests/mock_tests/Makefile.am index 7bbdda9ffba..5083924aa55 100644 --- a/tests/mock_tests/Makefile.am +++ b/tests/mock_tests/Makefile.am @@ -85,6 +85,7 @@ tests_SOURCES = aclorch_ut.cpp \ mock_orch_test.cpp \ mock_dash_orch_test.cpp \ zmq_orch_ut.cpp \ + zmq_route_orch_ut.cpp \ retrycache_ut.cpp \ mock_saihelper.cpp \ mirrororch_ut.cpp \ @@ -161,6 +162,7 @@ tests_SOURCES = aclorch_ut.cpp \ $(top_srcdir)/cfgmgr/portmgr.cpp \ $(top_srcdir)/cfgmgr/sflowmgr.cpp \ $(top_srcdir)/orchagent/zmqorch.cpp \ + $(top_srcdir)/orchagent/zmqrouteorch.cpp \ $(top_srcdir)/orchagent/dash/dashenifwdorch.cpp \ $(top_srcdir)/orchagent/dash/dashenifwdinfo.cpp \ $(top_srcdir)/orchagent/dash/dashaclorch.cpp \ diff --git a/tests/mock_tests/zmq_route_orch_ut.cpp b/tests/mock_tests/zmq_route_orch_ut.cpp new file mode 100644 index 00000000000..3d94c05cb40 --- /dev/null +++ b/tests/mock_tests/zmq_route_orch_ut.cpp @@ -0,0 +1,359 @@ +#include +#include +#include +#include +#include +#include + +#include "gtest/gtest.h" +#include "schema.h" +#include "ut_helper.h" +#include "orch_zmq_config.h" +#include "dbconnector.h" +#include "mock_table.h" +#include "select.h" +#include "zmqclient.h" +#include "zmqproducerstatetable.h" +#include "zmqrouteserver.h" +#include "zmqrouteconsumerstatetable.h" + +#define protected public +#include "orch.h" +#include "zmqrouteorch.h" +#undef protected + +using namespace std; +using namespace swss; + +extern size_t gMaxBulkSize; + +namespace { + +// Wait until pred() becomes true or deadlineMs elapses; returns the final value. +template +bool waitFor(int deadlineMs, Pred pred) +{ + auto deadline = std::chrono::steady_clock::now() + std::chrono::milliseconds(deadlineMs); + while (std::chrono::steady_clock::now() < deadline) + { + if (pred()) + return true; + std::this_thread::sleep_for(std::chrono::milliseconds(1)); + } + return pred(); +} + +// Minimal subclass of ZmqRouteOrch that records doTask invocations, so tests +// can assert that drain() forwards correctly without needing a full RouteOrch. +class RecordingZmqRouteOrch : public ZmqRouteOrch +{ +public: + RecordingZmqRouteOrch(swss::DBConnector *db, + const std::vector &tables, + ZmqRouteServer *zmqServer) + : ZmqRouteOrch(db, tables, zmqServer) + { + } + + void doTask(ConsumerBase &consumer) override + { + ++doTaskCount; + // Drain the consumer's m_toSync so subsequent drain() calls observe it + // as empty (matches the contract a real orch would honor). + consumer.m_toSync.clear(); + } + + std::atomic doTaskCount{0}; +}; + +} // namespace + +// ZmqRouteOrch with a nullptr server falls back to plain Consumer (legacy +// non-ZMQ path) for APPL_DB tables. +TEST(ZmqRouteOrchTest, NullServerFallsBackToConsumer) +{ + vector tables = { + { "ZMQ_ROUTE_UT_T1", 1 }, + { "ZMQ_ROUTE_UT_T2", 2 }, + }; + auto app_db = make_shared("APPL_DB", 0); + auto orch = make_shared(app_db.get(), tables, nullptr); + + EXPECT_EQ(orch->getSelectables().size(), tables.size()); + // Ensure the executor is a plain Consumer (not a ZmqRouteConsumer): the + // legacy fallback shouldn't pull in the ZmqRouteConsumer machinery. + auto exec = orch->m_consumerMap.begin()->second.get(); + EXPECT_EQ(dynamic_cast(exec), nullptr); +} + +// vector ctor (no per-table priority) — exercises the +// default_orch_pri code path in ZmqRouteOrch::ZmqRouteOrch(vector,...). +TEST(ZmqRouteOrchTest, VectorOfStringsCtor) +{ + vector tables = { "ZMQ_ROUTE_UT_TS1", "ZMQ_ROUTE_UT_TS2" }; + auto app_db = make_shared("APPL_DB", 0); + auto orch = make_shared(app_db.get(), tables, nullptr); + EXPECT_EQ(orch->getSelectables().size(), tables.size()); +} + +// Non-APPL_DB databases are unsupported; addConsumer should warn and create +// no executor. +TEST(ZmqRouteOrchTest, UnsupportedDbProducesNoExecutor) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto state_db = make_shared("STATE_DB", 0); + auto orch = make_shared(state_db.get(), tables, nullptr); + EXPECT_EQ(orch->getSelectables().size(), 0u); +} + +// With a real ZmqRouteServer, ZmqRouteOrch creates a ZmqRouteConsumer (not a +// plain Consumer). The server must outlive the orch. +TEST(ZmqRouteOrchTest, RealServerCreatesZmqRouteConsumer) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1260", "", /*lazyBind=*/true); + + auto orch = make_shared(app_db.get(), tables, &server); + ASSERT_EQ(orch->getSelectables().size(), tables.size()); + + auto exec = orch->m_consumerMap.begin()->second.get(); + EXPECT_NE(dynamic_cast(exec), nullptr); +} + +// doTask(Consumer&) on the base ZmqRouteOrch is a stub that forwards to the +// virtual doTask(ConsumerBase&) — this is the only piece that ZmqRouteOrch +// itself implements (besides ctors / addConsumer). Cover it. +TEST(ZmqRouteOrchTest, DoTaskConsumerForwardsToConsumerBase) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + auto orch = make_shared(app_db.get(), tables, nullptr); + + auto *exec = orch->m_consumerMap.begin()->second.get(); + auto *consumer = dynamic_cast(exec); + ASSERT_NE(consumer, nullptr); + + // Forge a single entry into m_toSync so that the recording doTask can see + // something and so the subsequent clear() actually does work. SyncMap is a + // multimap, so use insert rather than operator[]. + consumer->m_toSync.insert({ + "k1", + std::make_tuple(std::string("k1"), std::string(SET_COMMAND), + std::vector{{"f", "v"}}) + }); + + // ZmqRouteOrch::doTask(Consumer&) forwards to doTask(ConsumerBase&). + static_cast(orch.get())->doTask(*consumer); + EXPECT_EQ(orch->doTaskCount.load(), 1); + EXPECT_TRUE(consumer->m_toSync.empty()); +} + +// Drain on a ZmqRouteConsumer with empty m_toSync must NOT call doTask. +// Drain on a non-empty m_toSync must call doTask exactly once and the lock +// must allow re-entry afterwards. +TEST(ZmqRouteConsumerTest, DrainGatedByToSyncEmptiness) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1261", "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + + auto *exec = orch->m_consumerMap.begin()->second.get(); + auto *zrc = dynamic_cast(exec); + ASSERT_NE(zrc, nullptr); + + // Empty m_toSync: drain is a no-op. + zrc->drain(); + EXPECT_EQ(orch->doTaskCount.load(), 0); + + // Stage one entry via the locked addToSync override; drain forwards to + // doTask exactly once. RecordingZmqRouteOrch::doTask clears m_toSync. + KeyOpFieldsValuesTuple kfv("route_a", SET_COMMAND, + vector{{"f", "v"}}); + zrc->addToSync(kfv); + EXPECT_EQ(zrc->m_toSync.size(), 1u); + + zrc->drain(); + EXPECT_EQ(orch->doTaskCount.load(), 1); + EXPECT_TRUE(zrc->m_toSync.empty()); + + // A subsequent empty drain still doesn't call doTask, and the lock + // re-acquires cleanly. + zrc->drain(); + EXPECT_EQ(orch->doTaskCount.load(), 1); +} + +// execute() simply calls drain(); cover that override. +TEST(ZmqRouteConsumerTest, ExecuteDelegatesToDrain) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1262", "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + KeyOpFieldsValuesTuple kfv("route_b", SET_COMMAND, + vector{{"f", "v"}}); + zrc->addToSync(kfv); + + zrc->execute(); + EXPECT_EQ(orch->doTaskCount.load(), 1); +} + +// Locked deque-form addToSync forwards to ConsumerBase::addToSync(deque) and +// returns the count. +TEST(ZmqRouteConsumerTest, AddToSyncDequeReturnsCount) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1263", "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + std::deque entries; + for (int i = 0; i < 5; ++i) + { + entries.emplace_back("k" + std::to_string(i), SET_COMMAND, + vector{{"f", "v"}}); + } + + EXPECT_EQ(zrc->addToSync(entries), 5u); + EXPECT_EQ(zrc->m_toSync.size(), 5u); +} + +// dumpPendingTasks (locked override) returns the staged entries as strings +// and doesn't deadlock with concurrent addToSync. +TEST(ZmqRouteConsumerTest, DumpPendingTasksLockedAndCorrect) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1264", "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + zrc->addToSync(KeyOpFieldsValuesTuple("kA", SET_COMMAND, + vector{{"f", "v"}})); + zrc->addToSync(KeyOpFieldsValuesTuple("kB", DEL_COMMAND, + vector{})); + + std::vector ts; + zrc->dumpPendingTasks(ts); + EXPECT_EQ(ts.size(), 2u); +} + +// Concurrent addToSync from multiple threads must not crash, lose entries, or +// deadlock with drain. This guards the locking contract that ZmqRouteServer +// relies on (mqPollThread races with the orch main thread). +TEST(ZmqRouteConsumerTest, ConcurrentAddToSyncIsThreadSafe) +{ + vector tables = { { "ZMQ_ROUTE_UT_T1", 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server("tcp://*:1265", "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + constexpr int kThreads = 4; + constexpr int kPerThread = 250; + std::vector producers; + for (int t = 0; t < kThreads; ++t) + { + producers.emplace_back([zrc, t]() { + for (int i = 0; i < kPerThread; ++i) + { + std::string k = "t" + std::to_string(t) + "_" + std::to_string(i); + zrc->addToSync(KeyOpFieldsValuesTuple( + k, SET_COMMAND, vector{{"f", "v"}})); + } + }); + } + for (auto &th : producers) + th.join(); + + EXPECT_EQ(zrc->m_toSync.size(), + static_cast(kThreads * kPerThread)); +} + +// End-to-end: ZmqProducerStateTable → ZmqRouteServer → ZmqRouteConsumer +// ingress callback → m_toSync. Verifies the callback wiring set up by +// ZmqRouteConsumer's constructor actually merges entries into m_toSync, and +// (since count < gMaxBulkSize) does not eagerly fire notifyPending — the +// burst quiesce timer fires it instead. +TEST(ZmqRouteConsumerTest, IngressCallbackMergesIntoToSync) +{ + const string tableName = "ZMQ_ROUTE_UT_INGRESS"; + const string pushEndpoint = "tcp://localhost:1266"; + const string pullEndpoint = "tcp://*:1266"; + + vector tables = { { tableName, 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server(pullEndpoint, "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + server.bind(); + + ZmqClient client(pushEndpoint, 0); + ZmqProducerStateTable p(app_db.get(), tableName, client, /*dbPersistence=*/false); + p.set("route_x", vector{{"nh", "1.1.1.1"}}); + + ASSERT_TRUE(waitFor(2000, [&] { return zrc->m_toSync.size() >= 1u; })); + EXPECT_NE(zrc->m_toSync.find("route_x"), zrc->m_toSync.end()); +} + +// When the ingress callback fills m_toSync past gMaxBulkSize, it must fire +// notifyPending mid-burst (rather than waiting for the burst quiesce timer) +// so the orch main loop wakes up and drains immediately. We lower +// gMaxBulkSize to 1 to make this trivially observable. +TEST(ZmqRouteConsumerTest, IngressCallbackFiresNotifyAtMaxBulkSize) +{ + const string tableName = "ZMQ_ROUTE_UT_BULK"; + const string pushEndpoint = "tcp://localhost:1267"; + const string pullEndpoint = "tcp://*:1267"; + + vector tables = { { tableName, 1 } }; + auto app_db = make_shared("APPL_DB", 0); + ZmqRouteServer server(pullEndpoint, "", /*lazyBind=*/true); + auto orch = make_shared(app_db.get(), tables, &server); + auto *zrc = dynamic_cast( + orch->m_consumerMap.begin()->second.get()); + ASSERT_NE(zrc, nullptr); + + server.bind(); + + // Force the mid-burst notify branch to trip on the very first callback. + const size_t savedMaxBulk = gMaxBulkSize; + gMaxBulkSize = 1; + + ZmqClient client(pushEndpoint, 0); + ZmqProducerStateTable p(app_db.get(), tableName, client, /*dbPersistence=*/false); + p.set("route_bulk", vector{{"nh", "2.2.2.2"}}); + + ASSERT_TRUE(waitFor(2000, [&] { return zrc->m_toSync.size() >= 1u; })); + + // Select wake-up should arrive almost immediately because the ingress + // callback fires notifyPending the moment m_toSync reaches gMaxBulkSize=1 + // — without this we'd have to wait for BURST_QUIESCE_MS (~5ms) before the + // post-burst notify fires. + Select sel; + sel.addSelectable(zrc); + Selectable *out = nullptr; + EXPECT_EQ(sel.select(&out, 200), Select::OBJECT); + EXPECT_EQ(out, zrc); + + gMaxBulkSize = savedMaxBulk; +}