From a49ab5ba256dbfbbebccd501dd03086aea2662b1 Mon Sep 17 00:00:00 2001 From: Martino Ferrari Date: Mon, 10 Aug 2026 17:18:50 +0200 Subject: [PATCH] Added silence timeout as floating point --- .../UDPStreamerClient/UDPStreamerClient.cpp | 11 ++++ .../UDPStreamerClient/UDPStreamerClient.h | 1 + .../Interfaces/UDPStream/UDPSClient.cpp | 6 +- .../Interfaces/UDPStream/UDPSClient.h | 7 ++- .../UDPStreamerClientGTest.cpp | 5 ++ .../UDPStreamerClientTest.cpp | 45 ++++++++++++++ .../UDPStreamerClient/UDPStreamerClientTest.h | 6 ++ Test/GTest/UDPSClientGTest.cpp | 60 +++++++++++++++++-- 8 files changed, 132 insertions(+), 9 deletions(-) diff --git a/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.cpp b/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.cpp index c15ce88..f93feef 100644 --- a/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.cpp +++ b/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.cpp @@ -58,6 +58,9 @@ static const uint32 UDPS_CLIENT_DEFAULT_MAX_PAYLOAD = 1400u; /** Default unicast keepalive interval (seconds); 0 disables. */ static const uint32 UDPS_CLIENT_DEFAULT_KEEPALIVE_INTERVAL_S = 15u; +/** Default silence timeout before reconnect (seconds); sub-second values allowed. */ +static const float32 UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S = 1.0f; + /** Bytes prepended to each DATA payload for the HRT packet timestamp. */ static const uint32 UDPS_CLIENT_TIMESTAMP_BYTES = 8u; @@ -133,6 +136,7 @@ UDPStreamerClient::UDPStreamerClient() : port = UDPS_CLIENT_DEFAULT_PORT; maxPayloadSize = UDPS_CLIENT_DEFAULT_MAX_PAYLOAD; keepAliveInterval = UDPS_CLIENT_DEFAULT_KEEPALIVE_INTERVAL_S; + silenceTimeout = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S; cpuMask = 0xFFFFFFFFu; stackSize = THREADS_DEFAULT_STACKSIZE; dataPort = UDPS_CLIENT_DEFAULT_PORT + UDPS_CLIENT_DEFAULT_DP_OFFSET; @@ -211,6 +215,12 @@ bool UDPStreamerClient::Initialise(StructuredDataI &data) { } } + if (ok) { + if (!data.Read("SilenceTimeout", silenceTimeout)) { + silenceTimeout = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S; + } + } + if (ok) { if (!data.Read("CPUMask", cpuMask)) { cpuMask = 0xFFFFFFFFu; @@ -256,6 +266,7 @@ bool UDPStreamerClient::Initialise(StructuredDataI &data) { } if (ok) { ok = cdb.Write("MaxPayloadSize", maxPayloadSize); } if (ok) { ok = cdb.Write("KeepAliveInterval", keepAliveInterval); } + if (ok) { ok = cdb.Write("SilenceTimeout", silenceTimeout); } if (ok) { ok = cdb.Write("CPUMask", cpuMask); } if (ok) { ok = cdb.Write("StackSize", stackSize); } if (ok) { ok = cdb.MoveToRoot(); } diff --git a/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.h b/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.h index 13c6177..26be648 100644 --- a/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.h +++ b/Source/Components/DataSources/UDPStreamerClient/UDPStreamerClient.h @@ -175,6 +175,7 @@ private: uint16 port; /**< Server port. */ uint32 maxPayloadSize; /**< Max payload bytes per datagram. */ uint32 keepAliveInterval; /**< Seconds between unicast keepalive ACKs (0 disables). */ + float32 silenceTimeout; /**< Seconds of no data before reconnect (sub-second allowed, 0 disables). */ uint32 cpuMask; /**< Background thread CPU affinity. */ uint32 stackSize; /**< Background thread stack size. */ StreamString multicastGroup; /**< Multicast group IP; empty = unicast. */ diff --git a/Source/Components/Interfaces/UDPStream/UDPSClient.cpp b/Source/Components/Interfaces/UDPStream/UDPSClient.cpp index 4a0887d..8c9ed05 100644 --- a/Source/Components/Interfaces/UDPStream/UDPSClient.cpp +++ b/Source/Components/Interfaces/UDPStream/UDPSClient.cpp @@ -87,9 +87,11 @@ bool UDPSClient::Initialise(StructuredDataI &data) { dataPort = static_cast(dpU32); } - uint32 silenceS = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S; + float32 silenceS = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S; (void) data.Read("SilenceTimeout", silenceS); - silenceTimeoutTicks = static_cast(silenceS) * HighResolutionTimer::Frequency(); + /* float64 math: the tick rate (~1e9) exceeds float32's 24-bit mantissa */ + silenceTimeoutTicks = static_cast(static_cast(silenceS) * + static_cast(HighResolutionTimer::Frequency())); uint32 reconnectS = UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S; (void) data.Read("ReconnectDelay", reconnectS); diff --git a/Source/Components/Interfaces/UDPStream/UDPSClient.h b/Source/Components/Interfaces/UDPStream/UDPSClient.h index 6439b96..7eec743 100644 --- a/Source/Components/Interfaces/UDPStream/UDPSClient.h +++ b/Source/Components/Interfaces/UDPStream/UDPSClient.h @@ -86,8 +86,8 @@ public: * 256-fragment span the recvMask[32] tracks at typical chunk sizes. */ static const uint32 UDPS_CLIENT_MAX_PACKET_BYTES = 1048576u; // 1 MiB - /** Default silence timeout before reconnect (seconds). */ - static const uint32 UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S = 5u; + /** Default silence timeout before reconnect (seconds); sub-second values allowed. */ + static const float32 UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S = 1.0f; /** Default delay between reconnect attempts (seconds). */ static const uint32 UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S = 2u; @@ -118,7 +118,8 @@ public: * - MulticastGroup (char*) IPv4 multicast address; presence enables multicast mode. * - Interface (char*) Network interface for multicast join (e.g. "lo"). Required when MulticastGroup is set. * - DataPort (uint16) UDP multicast data port (defaults to Port+1). - * - SilenceTimeout (uint32) Seconds of no data before reconnect. Default 5. + * - SilenceTimeout (float32) Seconds of no data before reconnect. Default 1.0. + * Sub-second values allowed; 0 disables the check. * - ReconnectDelay (uint32) Seconds to wait between reconnect attempts. Default 2. * - KeepAliveInterval (uint32) Seconds between unicast keepalive ACKs. Default 15. 0 disables. * - MaxPayloadSize (uint32) Max payload bytes per datagram, excluding header. Default 1400. diff --git a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientGTest.cpp b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientGTest.cpp index c94c6e1..6ae2b37 100644 --- a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientGTest.cpp +++ b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientGTest.cpp @@ -51,6 +51,11 @@ TEST(UDPStreamerClientGTest, TestInitialise_DefaultPort) { ASSERT_TRUE(test.TestInitialise_DefaultPort()); } +TEST(UDPStreamerClientGTest, TestInitialise_SilenceTimeoutFloat) { + UDPStreamerClientTest test; + ASSERT_TRUE(test.TestInitialise_SilenceTimeoutFloat()); +} + TEST(UDPStreamerClientGTest, TestInitialise_MulticastMode_Valid) { UDPStreamerClientTest test; ASSERT_TRUE(test.TestInitialise_MulticastMode_Valid()); diff --git a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.cpp b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.cpp index db1ae4f..2e5dfd9 100644 --- a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.cpp +++ b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.cpp @@ -354,6 +354,51 @@ bool UDPStreamerClientTest::TestInitialise_DefaultPort() { return ok; } +bool UDPStreamerClientTest::TestInitialise_SilenceTimeoutFloat() { + using namespace MARTe; + static const char8 *const cfg = + "+Test = {\n" + " Class = RealTimeApplication\n" + " +Functions = {\n" + " Class = ReferenceContainer\n" + " +Reader = {\n" + " Class = UDPStreamerClientTestGAM\n" + " InputSignals = { Counter = { DataSource = ClientDS Type = uint32 } }\n" + " OutputSignals = { Counter = { DataSource = DDB Type = uint32 } }\n" + " }\n" + " }\n" + " +Data = {\n" + " Class = ReferenceContainer\n" + " DefaultDataSource = DDB\n" + " +DDB = { Class = GAMDataSource }\n" + " +ClientDS = {\n" + " Class = UDPStreamerClient\n" + " ServerAddress = \"127.0.0.1\"\n" + " Port = 44702\n" + " SilenceTimeout = 0.25\n" + " Signals = { Counter = { Type = uint32 } }\n" + " }\n" + " +Timings = { Class = TimingDataSource }\n" + " }\n" + " +States = {\n" + " Class = ReferenceContainer\n" + " +State1 = {\n" + " Class = RealTimeState\n" + " +Threads = {\n" + " Class = ReferenceContainer\n" + " +Thread1 = { Class = RealTimeThread Functions = { Reader } }\n" + " }\n" + " }\n" + " }\n" + " +Scheduler = { Class = GAMScheduler TimingDataSource = Timings }\n" + "}\n"; + + ReferenceT app = LoadApplication(cfg); + bool ok = app.IsValid(); + ObjectRegistryDatabase::Instance()->Purge(); + return ok; +} + bool UDPStreamerClientTest::TestInitialise_MulticastMode_Valid() { using namespace MARTe; static const char8 *const cfg = diff --git a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.h b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.h index a6ce9e5..f4c0f87 100644 --- a/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.h +++ b/Test/Components/DataSources/UDPStreamerClient/UDPStreamerClientTest.h @@ -62,6 +62,12 @@ public: */ bool TestInitialise_MulticastMode_Valid(); + /** + * @brief Tests Initialise with a sub-second float32 SilenceTimeout (0.25 s) + * forwarded through to the UDPSClient receiver. + */ + bool TestInitialise_SilenceTimeoutFloat(); + /** * @brief Tests that DataPort defaults to Port+1 when MulticastGroup is * set but DataPort is absent. diff --git a/Test/GTest/UDPSClientGTest.cpp b/Test/GTest/UDPSClientGTest.cpp index 8cc3ab7..a865cd6 100644 --- a/Test/GTest/UDPSClientGTest.cpp +++ b/Test/GTest/UDPSClientGTest.cpp @@ -150,7 +150,7 @@ TEST(UDPSClientGTest, TestUnicastKeepAliveSendsPeriodicAck) { ASSERT_TRUE(cfg.Write("Port", static_cast(serverPort))); ASSERT_TRUE(cfg.Write("KeepAliveInterval", 1u)); /* SilenceTimeout=0 keeps the session stable for the whole test */ - ASSERT_TRUE(cfg.Write("SilenceTimeout", 0u)); + ASSERT_TRUE(cfg.Write("SilenceTimeout", 0.0f)); UDPSClient client; ASSERT_TRUE(client.Initialise(cfg)); @@ -195,7 +195,7 @@ TEST(UDPSClientGTest, TestKeepAliveDisabledWhenIntervalZero) { ASSERT_TRUE(cfg.Write("ServerAddr", "127.0.0.1")); ASSERT_TRUE(cfg.Write("Port", static_cast(serverPort))); ASSERT_TRUE(cfg.Write("KeepAliveInterval", 0u)); - ASSERT_TRUE(cfg.Write("SilenceTimeout", 0u)); + ASSERT_TRUE(cfg.Write("SilenceTimeout", 0.0f)); UDPSClient client; ASSERT_TRUE(client.Initialise(cfg)); @@ -240,7 +240,7 @@ TEST(UDPSClientGTest, TestKeepAlivePreventsServerEviction) { ASSERT_TRUE(clientCfg.Write("ServerAddr", "127.0.0.1")); ASSERT_TRUE(clientCfg.Write("Port", static_cast(serverPort))); ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 1u)); - ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0u)); + ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0.0f)); UDPSClient client; ASSERT_TRUE(client.Initialise(clientCfg)); ASSERT_TRUE(client.Start()); @@ -279,7 +279,7 @@ TEST(UDPSClientGTest, TestServerEvictsWithoutKeepAlive) { ASSERT_TRUE(clientCfg.Write("ServerAddr", "127.0.0.1")); ASSERT_TRUE(clientCfg.Write("Port", static_cast(serverPort))); ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 0u)); - ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0u)); + ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0.0f)); UDPSClient client; ASSERT_TRUE(client.Initialise(clientCfg)); ASSERT_TRUE(client.Start()); @@ -294,3 +294,55 @@ TEST(UDPSClientGTest, TestServerEvictsWithoutKeepAlive) { client.Stop(); server.Stop(); } + +TEST(UDPSClientGTest, TestSilenceTimeoutSubSecondTriggersReconnect) { + /* SilenceTimeout is float32 seconds: a sub-second value must actually + * fire (integer truncation would silently disable the check). */ + BasicUDPSocket server; + ASSERT_TRUE(server.Open()); + ASSERT_TRUE(server.Listen(0u)); + uint16 serverPort = GetBoundPort(server); + ASSERT_NE(serverPort, 0u); + + ConfigurationDatabase cfg; + ASSERT_TRUE(cfg.Write("ServerAddr", "127.0.0.1")); + ASSERT_TRUE(cfg.Write("Port", static_cast(serverPort))); + ASSERT_TRUE(cfg.Write("SilenceTimeout", 0.3f)); + ASSERT_TRUE(cfg.Write("ReconnectDelay", 0u)); /* reconnect immediately */ + ASSERT_TRUE(cfg.Write("KeepAliveInterval", 0u)); + + UDPSClient client; + ASSERT_TRUE(client.Initialise(cfg)); + ASSERT_TRUE(client.Start()); + + /* 1) First CONNECT from the client's ephemeral socket */ + uint8 type = 0xFFu; + uint16 portA = 0u; + ASSERT_TRUE(WaitDatagram(server, 3000, type, portA)); + EXPECT_EQ(type, UDPS_TYPE_CONNECT); + ASSERT_NE(portA, 0u); + + /* 2) The server sends nothing: after ~0.3 s the client must disconnect + * and re-announce with a NEW ephemeral socket. Fails if the timeout + * was truncated to 0 (disabled) or left at the old 5 s default. */ + uint32 elapsedMs = 0u; + bool reconnected = false; + while ((elapsedMs < 2000u) && !reconnected) { + uint8 t = 0xFFu; + uint16 p = 0u; + bool got = WaitDatagram(server, 500, t, p); + elapsedMs += 500u; + if (!got) { + continue; + } + /* DISCONNECT from the old socket is expected; only a CONNECT from a + * new source port proves the reconnect happened. */ + if ((t == UDPS_TYPE_CONNECT) && (p != portA)) { + reconnected = true; + } + } + EXPECT_TRUE(reconnected); + + client.Stop(); + server.Close(); +}