Added silence timeout as floating point
This commit is contained in:
@@ -58,6 +58,9 @@ static const uint32 UDPS_CLIENT_DEFAULT_MAX_PAYLOAD = 1400u;
|
|||||||
/** Default unicast keepalive interval (seconds); 0 disables. */
|
/** Default unicast keepalive interval (seconds); 0 disables. */
|
||||||
static const uint32 UDPS_CLIENT_DEFAULT_KEEPALIVE_INTERVAL_S = 15u;
|
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. */
|
/** Bytes prepended to each DATA payload for the HRT packet timestamp. */
|
||||||
static const uint32 UDPS_CLIENT_TIMESTAMP_BYTES = 8u;
|
static const uint32 UDPS_CLIENT_TIMESTAMP_BYTES = 8u;
|
||||||
|
|
||||||
@@ -133,6 +136,7 @@ UDPStreamerClient::UDPStreamerClient() :
|
|||||||
port = UDPS_CLIENT_DEFAULT_PORT;
|
port = UDPS_CLIENT_DEFAULT_PORT;
|
||||||
maxPayloadSize = UDPS_CLIENT_DEFAULT_MAX_PAYLOAD;
|
maxPayloadSize = UDPS_CLIENT_DEFAULT_MAX_PAYLOAD;
|
||||||
keepAliveInterval = UDPS_CLIENT_DEFAULT_KEEPALIVE_INTERVAL_S;
|
keepAliveInterval = UDPS_CLIENT_DEFAULT_KEEPALIVE_INTERVAL_S;
|
||||||
|
silenceTimeout = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S;
|
||||||
cpuMask = 0xFFFFFFFFu;
|
cpuMask = 0xFFFFFFFFu;
|
||||||
stackSize = THREADS_DEFAULT_STACKSIZE;
|
stackSize = THREADS_DEFAULT_STACKSIZE;
|
||||||
dataPort = UDPS_CLIENT_DEFAULT_PORT + UDPS_CLIENT_DEFAULT_DP_OFFSET;
|
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 (ok) {
|
||||||
if (!data.Read("CPUMask", cpuMask)) {
|
if (!data.Read("CPUMask", cpuMask)) {
|
||||||
cpuMask = 0xFFFFFFFFu;
|
cpuMask = 0xFFFFFFFFu;
|
||||||
@@ -256,6 +266,7 @@ bool UDPStreamerClient::Initialise(StructuredDataI &data) {
|
|||||||
}
|
}
|
||||||
if (ok) { ok = cdb.Write("MaxPayloadSize", maxPayloadSize); }
|
if (ok) { ok = cdb.Write("MaxPayloadSize", maxPayloadSize); }
|
||||||
if (ok) { ok = cdb.Write("KeepAliveInterval", keepAliveInterval); }
|
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("CPUMask", cpuMask); }
|
||||||
if (ok) { ok = cdb.Write("StackSize", stackSize); }
|
if (ok) { ok = cdb.Write("StackSize", stackSize); }
|
||||||
if (ok) { ok = cdb.MoveToRoot(); }
|
if (ok) { ok = cdb.MoveToRoot(); }
|
||||||
|
|||||||
@@ -175,6 +175,7 @@ private:
|
|||||||
uint16 port; /**< Server port. */
|
uint16 port; /**< Server port. */
|
||||||
uint32 maxPayloadSize; /**< Max payload bytes per datagram. */
|
uint32 maxPayloadSize; /**< Max payload bytes per datagram. */
|
||||||
uint32 keepAliveInterval; /**< Seconds between unicast keepalive ACKs (0 disables). */
|
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 cpuMask; /**< Background thread CPU affinity. */
|
||||||
uint32 stackSize; /**< Background thread stack size. */
|
uint32 stackSize; /**< Background thread stack size. */
|
||||||
StreamString multicastGroup; /**< Multicast group IP; empty = unicast. */
|
StreamString multicastGroup; /**< Multicast group IP; empty = unicast. */
|
||||||
|
|||||||
@@ -87,9 +87,11 @@ bool UDPSClient::Initialise(StructuredDataI &data) {
|
|||||||
dataPort = static_cast<uint16>(dpU32);
|
dataPort = static_cast<uint16>(dpU32);
|
||||||
}
|
}
|
||||||
|
|
||||||
uint32 silenceS = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S;
|
float32 silenceS = UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S;
|
||||||
(void) data.Read("SilenceTimeout", silenceS);
|
(void) data.Read("SilenceTimeout", silenceS);
|
||||||
silenceTimeoutTicks = static_cast<uint64>(silenceS) * HighResolutionTimer::Frequency();
|
/* float64 math: the tick rate (~1e9) exceeds float32's 24-bit mantissa */
|
||||||
|
silenceTimeoutTicks = static_cast<uint64>(static_cast<float64>(silenceS) *
|
||||||
|
static_cast<float64>(HighResolutionTimer::Frequency()));
|
||||||
|
|
||||||
uint32 reconnectS = UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S;
|
uint32 reconnectS = UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S;
|
||||||
(void) data.Read("ReconnectDelay", reconnectS);
|
(void) data.Read("ReconnectDelay", reconnectS);
|
||||||
|
|||||||
@@ -86,8 +86,8 @@ public:
|
|||||||
* 256-fragment span the recvMask[32] tracks at typical chunk sizes. */
|
* 256-fragment span the recvMask[32] tracks at typical chunk sizes. */
|
||||||
static const uint32 UDPS_CLIENT_MAX_PACKET_BYTES = 1048576u; // 1 MiB
|
static const uint32 UDPS_CLIENT_MAX_PACKET_BYTES = 1048576u; // 1 MiB
|
||||||
|
|
||||||
/** Default silence timeout before reconnect (seconds). */
|
/** Default silence timeout before reconnect (seconds); sub-second values allowed. */
|
||||||
static const uint32 UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S = 5u;
|
static const float32 UDPS_CLIENT_DEFAULT_SILENCE_TIMEOUT_S = 1.0f;
|
||||||
|
|
||||||
/** Default delay between reconnect attempts (seconds). */
|
/** Default delay between reconnect attempts (seconds). */
|
||||||
static const uint32 UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S = 2u;
|
static const uint32 UDPS_CLIENT_DEFAULT_RECONNECT_DELAY_S = 2u;
|
||||||
@@ -118,7 +118,8 @@ public:
|
|||||||
* - MulticastGroup (char*) IPv4 multicast address; presence enables multicast mode.
|
* - MulticastGroup (char*) IPv4 multicast address; presence enables multicast mode.
|
||||||
* - Interface (char*) Network interface for multicast join (e.g. "lo"). Required when MulticastGroup is set.
|
* - 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).
|
* - 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.
|
* - ReconnectDelay (uint32) Seconds to wait between reconnect attempts. Default 2.
|
||||||
* - KeepAliveInterval (uint32) Seconds between unicast keepalive ACKs. Default 15. 0 disables.
|
* - KeepAliveInterval (uint32) Seconds between unicast keepalive ACKs. Default 15. 0 disables.
|
||||||
* - MaxPayloadSize (uint32) Max payload bytes per datagram, excluding header. Default 1400.
|
* - MaxPayloadSize (uint32) Max payload bytes per datagram, excluding header. Default 1400.
|
||||||
|
|||||||
@@ -51,6 +51,11 @@ TEST(UDPStreamerClientGTest, TestInitialise_DefaultPort) {
|
|||||||
ASSERT_TRUE(test.TestInitialise_DefaultPort());
|
ASSERT_TRUE(test.TestInitialise_DefaultPort());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
TEST(UDPStreamerClientGTest, TestInitialise_SilenceTimeoutFloat) {
|
||||||
|
UDPStreamerClientTest test;
|
||||||
|
ASSERT_TRUE(test.TestInitialise_SilenceTimeoutFloat());
|
||||||
|
}
|
||||||
|
|
||||||
TEST(UDPStreamerClientGTest, TestInitialise_MulticastMode_Valid) {
|
TEST(UDPStreamerClientGTest, TestInitialise_MulticastMode_Valid) {
|
||||||
UDPStreamerClientTest test;
|
UDPStreamerClientTest test;
|
||||||
ASSERT_TRUE(test.TestInitialise_MulticastMode_Valid());
|
ASSERT_TRUE(test.TestInitialise_MulticastMode_Valid());
|
||||||
|
|||||||
@@ -354,6 +354,51 @@ bool UDPStreamerClientTest::TestInitialise_DefaultPort() {
|
|||||||
return ok;
|
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<RealTimeApplication> app = LoadApplication(cfg);
|
||||||
|
bool ok = app.IsValid();
|
||||||
|
ObjectRegistryDatabase::Instance()->Purge();
|
||||||
|
return ok;
|
||||||
|
}
|
||||||
|
|
||||||
bool UDPStreamerClientTest::TestInitialise_MulticastMode_Valid() {
|
bool UDPStreamerClientTest::TestInitialise_MulticastMode_Valid() {
|
||||||
using namespace MARTe;
|
using namespace MARTe;
|
||||||
static const char8 *const cfg =
|
static const char8 *const cfg =
|
||||||
|
|||||||
@@ -62,6 +62,12 @@ public:
|
|||||||
*/
|
*/
|
||||||
bool TestInitialise_MulticastMode_Valid();
|
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
|
* @brief Tests that DataPort defaults to Port+1 when MulticastGroup is
|
||||||
* set but DataPort is absent.
|
* set but DataPort is absent.
|
||||||
|
|||||||
@@ -150,7 +150,7 @@ TEST(UDPSClientGTest, TestUnicastKeepAliveSendsPeriodicAck) {
|
|||||||
ASSERT_TRUE(cfg.Write("Port", static_cast<uint32>(serverPort)));
|
ASSERT_TRUE(cfg.Write("Port", static_cast<uint32>(serverPort)));
|
||||||
ASSERT_TRUE(cfg.Write("KeepAliveInterval", 1u));
|
ASSERT_TRUE(cfg.Write("KeepAliveInterval", 1u));
|
||||||
/* SilenceTimeout=0 keeps the session stable for the whole test */
|
/* 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;
|
UDPSClient client;
|
||||||
ASSERT_TRUE(client.Initialise(cfg));
|
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("ServerAddr", "127.0.0.1"));
|
||||||
ASSERT_TRUE(cfg.Write("Port", static_cast<uint32>(serverPort)));
|
ASSERT_TRUE(cfg.Write("Port", static_cast<uint32>(serverPort)));
|
||||||
ASSERT_TRUE(cfg.Write("KeepAliveInterval", 0u));
|
ASSERT_TRUE(cfg.Write("KeepAliveInterval", 0u));
|
||||||
ASSERT_TRUE(cfg.Write("SilenceTimeout", 0u));
|
ASSERT_TRUE(cfg.Write("SilenceTimeout", 0.0f));
|
||||||
|
|
||||||
UDPSClient client;
|
UDPSClient client;
|
||||||
ASSERT_TRUE(client.Initialise(cfg));
|
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("ServerAddr", "127.0.0.1"));
|
||||||
ASSERT_TRUE(clientCfg.Write("Port", static_cast<uint32>(serverPort)));
|
ASSERT_TRUE(clientCfg.Write("Port", static_cast<uint32>(serverPort)));
|
||||||
ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 1u));
|
ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 1u));
|
||||||
ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0u));
|
ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0.0f));
|
||||||
UDPSClient client;
|
UDPSClient client;
|
||||||
ASSERT_TRUE(client.Initialise(clientCfg));
|
ASSERT_TRUE(client.Initialise(clientCfg));
|
||||||
ASSERT_TRUE(client.Start());
|
ASSERT_TRUE(client.Start());
|
||||||
@@ -279,7 +279,7 @@ TEST(UDPSClientGTest, TestServerEvictsWithoutKeepAlive) {
|
|||||||
ASSERT_TRUE(clientCfg.Write("ServerAddr", "127.0.0.1"));
|
ASSERT_TRUE(clientCfg.Write("ServerAddr", "127.0.0.1"));
|
||||||
ASSERT_TRUE(clientCfg.Write("Port", static_cast<uint32>(serverPort)));
|
ASSERT_TRUE(clientCfg.Write("Port", static_cast<uint32>(serverPort)));
|
||||||
ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 0u));
|
ASSERT_TRUE(clientCfg.Write("KeepAliveInterval", 0u));
|
||||||
ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0u));
|
ASSERT_TRUE(clientCfg.Write("SilenceTimeout", 0.0f));
|
||||||
UDPSClient client;
|
UDPSClient client;
|
||||||
ASSERT_TRUE(client.Initialise(clientCfg));
|
ASSERT_TRUE(client.Initialise(clientCfg));
|
||||||
ASSERT_TRUE(client.Start());
|
ASSERT_TRUE(client.Start());
|
||||||
@@ -294,3 +294,55 @@ TEST(UDPSClientGTest, TestServerEvictsWithoutKeepAlive) {
|
|||||||
client.Stop();
|
client.Stop();
|
||||||
server.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<uint32>(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();
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user