minor source cleanup and annotation;

r05a06_dev
Bryan Biedenkapp 4 days ago
parent 1040ba35e8
commit 4774bb94af

@ -69,7 +69,7 @@
#define TAG_TRANSFER_ACT_LOG "TRNSLOG"
#define TAG_TRANSFER_DIAG_LOG "TRNSDIAG"
#define TAG_TRANSFER_STATUS "TRNSSTS"
#define TAG_TRANSFER_PATCH_STATUS "TRNSPTCH"
#define TAG_TRANSFER_PATCH_STAT "TRNSPTCH"
#define TAG_ANNOUNCE "ANNC"
#define TAG_PEER_REPLICA "REPL"

@ -22,6 +22,10 @@ constexpr uint32_t PatchStatusRegistry::MIN_TTL_SECONDS;
constexpr uint32_t PatchStatusRegistry::MAX_TTL_SECONDS;
constexpr uint32_t PatchStatusRegistry::MAX_WAIT_MS;
// ---------------------------------------------------------------------------
// Public Class Members
// ---------------------------------------------------------------------------
/* Initializes a new instance of the PatchStatusRegistry class. */
PatchStatusRegistry::PatchStatusRegistry() :
@ -272,6 +276,10 @@ uint32_t PatchStatusRegistry::maxTtlSeconds() const
return m_maxTtlSeconds;
}
// ---------------------------------------------------------------------------
// Private Class Members
// ---------------------------------------------------------------------------
/* Gets the current system time in milliseconds. */
uint64_t PatchStatusRegistry::nowMs()

@ -24,6 +24,10 @@
#include <string>
#include <vector>
// ---------------------------------------------------------------------------
// Class Declaration
// ---------------------------------------------------------------------------
/**
* @brief In-memory registry for console-advertised patch status.
*

@ -156,31 +156,6 @@ void MetadataNetwork::close()
m_status = NET_STAT_INVALID;
}
/* Helper to send a metadata message to a peer's metadata port. */
bool MetadataNetwork::writePeerMetadata(FNEPeerConnection* connection, uint32_t ssrc, FrameQueue::OpcodePair opcode, const uint8_t* data,
uint32_t length, uint16_t pktSeq, uint32_t streamId) const
{
if (connection == nullptr)
return false;
if (m_status != NET_STAT_MST_RUNNING)
return false;
if (m_frameQueue == nullptr)
return false;
sockaddr_storage addr;
uint32_t addrLen = 0U;
uint16_t port = connection->port() + 1U;
if (udp::Socket::lookup(connection->address(), port, addr, addrLen) != 0) {
LogWarning(LOG_NET, "PEER %u (%s) failed to resolve metadata endpoint %s:%u", connection->id(),
connection->identWithQualifier().c_str(), connection->address().c_str(), port);
return false;
}
return m_frameQueue->write(data, length, streamId, connection->id(), ssrc, opcode, pktSeq, addr, addrLen);
}
// ---------------------------------------------------------------------------
// Private Class Members
// ---------------------------------------------------------------------------
@ -238,6 +213,31 @@ void MetadataNetwork::erasePacketBufferEntry(PacketBufferMap& pktMap, uint32_t p
pktMap.unlock();
}
/* Helper to send a metadata message to a peer's metadata port. */
bool MetadataNetwork::writePeerMetadata(FNEPeerConnection* connection, uint32_t ssrc, FrameQueue::OpcodePair opcode, const uint8_t* data,
uint32_t length, uint16_t pktSeq, uint32_t streamId) const
{
if (connection == nullptr)
return false;
if (m_status != NET_STAT_MST_RUNNING)
return false;
if (m_frameQueue == nullptr)
return false;
sockaddr_storage addr;
uint32_t addrLen = 0U;
uint16_t port = connection->port() + 1U;
if (udp::Socket::lookup(connection->address(), port, addr, addrLen) != 0) {
LogWarning(LOG_NET, "PEER %u (%s) failed to resolve metadata endpoint %s:%u", connection->id(),
connection->identWithQualifier().c_str(), connection->address().c_str(), port);
return false;
}
return m_frameQueue->write(data, length, streamId, connection->id(), ssrc, opcode, pktSeq, addr, addrLen);
}
/*
** Packet Processing
*/

@ -107,8 +107,6 @@ namespace network
friend class ::FNETestHooks;
#endif
friend class TrafficNetwork;
bool writePeerMetadata(FNEPeerConnection* connection, uint32_t ssrc, FrameQueue::OpcodePair opcode, const uint8_t* data,
uint32_t length, uint16_t pktSeq, uint32_t streamId) const;
TrafficNetwork* m_trafficNetwork;
HostFNE* m_host;
@ -181,6 +179,20 @@ namespace network
*/
static void erasePacketBufferEntry(PacketBufferMap& pktMap, uint32_t peerId, const PacketBufferEntryPtr& pkt);
/**
* @brief Writes metadata for a peer connection.
* @param connection Pointer to the FNEPeerConnection instance.
* @param ssrc SSRC of the packet.
* @param opcode Opcode pair of the frame.
* @param data Pointer to the data buffer.
* @param length Length of the data buffer.
* @param pktSeq Packet sequence number.
* @param streamId Stream ID of the packet.
* @returns True if the metadata was successfully written, false otherwise.
*/
bool writePeerMetadata(FNEPeerConnection* connection, uint32_t ssrc, FrameQueue::OpcodePair opcode, const uint8_t* data,
uint32_t length, uint16_t pktSeq, uint32_t streamId) const;
/*
** Packet Processing
*/

@ -832,9 +832,10 @@ void TrafficNetwork::clock(uint32_t ms)
}
}
// cleanup possibly stale data calls
if (m_patchStatusEnabled && m_patchStatusRegistry.cleanupExpired() > 0U)
writePatchStatusToConsoles(m_patchStatusRegistry.snapshot());
// cleanup possibly stale data calls
m_tagDMR->packetData()->cleanupStale();
m_tagP25->packetData()->cleanupStale();
@ -1569,162 +1570,6 @@ json::object TrafficNetwork::fneConnObject(uint32_t peerId, FNEPeerConnection *c
return peerObj;
}
/* Helper to send patch status state to one console peer. */
bool TrafficNetwork::writePatchStatusToPeer(uint32_t peerId, json::object obj)
{
if (peerId == 0U)
return false;
if (!m_patchStatusEnabled)
return false;
bool ret = false;
m_peers.shared_lock();
auto it = std::find_if(m_peers.begin(), m_peers.end(), [&](PeerMapPair x) { return x.first == peerId; });
if (it != m_peers.end() && it->second != nullptr)
ret = writePatchStatusPayload(it->second, obj);
m_peers.shared_unlock();
return ret;
}
/* Helper to broadcast patch status state to connected console peers. */
void TrafficNetwork::writePatchStatusToConsoles(json::object obj, uint32_t exceptPeerId)
{
if (!m_patchStatusEnabled)
return;
m_peers.shared_lock();
if (m_peers.size() == 0U) {
m_peers.shared_unlock();
return;
}
for (auto peer : m_peers) {
if (peer.first == exceptPeerId)
continue;
if (peer.second == nullptr)
continue;
if (!peer.second->connected() || peer.second->peerClass() != PEER_CONN_CLASS_CONSOLE)
continue;
writePatchStatusPayload(peer.second, obj);
}
m_peers.shared_unlock();
}
/* Helper to replicate patch status state to neighboring FNE peers. */
void TrafficNetwork::replicatePatchStatus(json::object obj, uint32_t exceptPeerId)
{
if (!m_patchStatusEnabled)
return;
obj["type"].set<std::string>("publish");
if (m_host->m_peerNetworks.size() > 0U) {
for (auto peer : m_host->m_peerNetworks) {
if (peer.first == exceptPeerId)
continue;
if (peer.second != nullptr && peer.second->isEnabled() && peer.second->isReplica())
peer.second->writePatchStatus(obj);
}
}
m_peers.shared_lock();
for (auto peer : m_peers) {
if (peer.first == exceptPeerId)
continue;
if (peer.second == nullptr)
continue;
if (!peer.second->connected() || peer.second->peerClass() != PEER_CONN_CLASS_NEIGHBOR || !peer.second->isReplica())
continue;
writePatchStatusReplicationPayload(peer.second, obj);
}
m_peers.shared_unlock();
}
/* Helper to serialize and queue a patch status transfer payload. */
bool TrafficNetwork::writePatchStatusPayload(FNEPeerConnection* connection, json::object obj)
{
if (connection == nullptr)
return false;
if (!m_patchStatusEnabled)
return false;
if (m_host->m_mdNetwork == nullptr)
return false;
if (!connection->connected())
return false;
if (connection->peerClass() != PEER_CONN_CLASS_CONSOLE)
return false;
obj["type"].set<std::string>("registry");
json::value v = json::value(obj);
std::string payload = std::string(v.serialize());
uint32_t len = static_cast<uint32_t>(payload.length());
if ((len + 11U) > DATA_PACKET_LENGTH) {
LogError(LOG_MASTER, "PEER %u (%s) patch status registry payload too large, len = %u", connection->id(), connection->identWithQualifier().c_str(), len);
return false;
}
uint8_t buffer[DATA_PACKET_LENGTH];
::memset(buffer, 0x00U, DATA_PACKET_LENGTH);
::memcpy(buffer + 11U, payload.c_str(), len);
if (m_debug) {
LogDebug(LOG_MASTER, "PEER %u (%s) sending patch status registry, len = %u", connection->id(), connection->identWithQualifier().c_str(), len);
}
return m_host->m_mdNetwork->writePeerMetadata(connection, m_peerId,
{ NET_FUNC::TRANSFER, NET_SUBFUNC::TRANSFER_SUBFUNC_PATCH_STATUS }, buffer, len + 11U, RTP_END_OF_CALL_SEQ, createStreamId());
}
/* Helper to serialize and queue a patch status replication payload. */
bool TrafficNetwork::writePatchStatusReplicationPayload(FNEPeerConnection* connection, json::object obj)
{
if (connection == nullptr)
return false;
if (!m_patchStatusEnabled)
return false;
if (m_host->m_mdNetwork == nullptr)
return false;
if (!connection->connected())
return false;
if (connection->peerClass() != PEER_CONN_CLASS_NEIGHBOR || !connection->isReplica())
return false;
obj["type"].set<std::string>("publish");
json::value v = json::value(obj);
std::string json = std::string(v.serialize());
size_t len = json.length() + 9U;
DECLARE_CHAR_ARRAY(buffer, len);
::memcpy(buffer + 0U, TAG_PEER_REPLICA, 4U);
::snprintf(buffer + 8U, json.length() + 1U, "%s", json.c_str());
PacketBuffer pkt(true, "Peer Replication, Patch Status");
pkt.encode((uint8_t*)buffer, len);
uint32_t streamId = createStreamId();
LogInfoEx(LOG_REPL, "PEER %u (%s) Peer Replication, Patch Status, blocks %u, streamId = %u", connection->id(),
connection->identWithQualifier().c_str(), pkt.fragments.size(), streamId);
if (pkt.fragments.size() > 0U) {
for (auto frag : pkt.fragments) {
m_host->m_mdNetwork->writePeerMetadata(connection, m_peerId, { NET_FUNC::REPL, NET_SUBFUNC::REPL_PATCH_STATUS },
frag.second->data, FRAG_SIZE, RTP_END_OF_CALL_SEQ, streamId);
Thread::sleep(60U); // pace block transmission
}
}
pkt.clear();
return true;
}
/* Helper to reset a peer connection. */
bool TrafficNetwork::resetPeer(uint32_t peerId)
@ -2102,6 +1947,166 @@ void TrafficNetwork::taskMetadataUpdate(MetadataUpdateRequest* req)
}
}
/*
** Console Patch Registry
*/
/* Helper to send patch status state to one console peer. */
bool TrafficNetwork::writePatchStatusToPeer(uint32_t peerId, json::object obj)
{
if (peerId == 0U)
return false;
if (!m_patchStatusEnabled)
return false;
bool ret = false;
m_peers.shared_lock();
auto it = std::find_if(m_peers.begin(), m_peers.end(), [&](PeerMapPair x) { return x.first == peerId; });
if (it != m_peers.end() && it->second != nullptr)
ret = writePatchStatusPayload(it->second, obj);
m_peers.shared_unlock();
return ret;
}
/* Helper to broadcast patch status state to connected console peers. */
void TrafficNetwork::writePatchStatusToConsoles(json::object obj, uint32_t exceptPeerId)
{
if (!m_patchStatusEnabled)
return;
m_peers.shared_lock();
if (m_peers.size() == 0U) {
m_peers.shared_unlock();
return;
}
for (auto peer : m_peers) {
if (peer.first == exceptPeerId)
continue;
if (peer.second == nullptr)
continue;
if (!peer.second->connected() || peer.second->peerClass() != PEER_CONN_CLASS_CONSOLE)
continue;
writePatchStatusPayload(peer.second, obj);
}
m_peers.shared_unlock();
}
/* Helper to replicate patch status state to neighboring FNE peers. */
void TrafficNetwork::replicatePatchStatus(json::object obj, uint32_t exceptPeerId)
{
if (!m_patchStatusEnabled)
return;
obj["type"].set<std::string>("publish");
if (m_host->m_peerNetworks.size() > 0U) {
for (auto peer : m_host->m_peerNetworks) {
if (peer.first == exceptPeerId)
continue;
if (peer.second != nullptr && peer.second->isEnabled() && peer.second->isReplica())
peer.second->writePatchStatus(obj);
}
}
m_peers.shared_lock();
for (auto peer : m_peers) {
if (peer.first == exceptPeerId)
continue;
if (peer.second == nullptr)
continue;
if (!peer.second->connected() || peer.second->peerClass() != PEER_CONN_CLASS_NEIGHBOR || !peer.second->isReplica())
continue;
writePatchStatusReplicationPayload(peer.second, obj);
}
m_peers.shared_unlock();
}
/* Helper to serialize and queue a patch status transfer payload. */
bool TrafficNetwork::writePatchStatusPayload(FNEPeerConnection* connection, json::object obj)
{
if (connection == nullptr)
return false;
if (!m_patchStatusEnabled)
return false;
if (m_host->m_mdNetwork == nullptr)
return false;
if (!connection->connected())
return false;
if (connection->peerClass() != PEER_CONN_CLASS_CONSOLE)
return false;
obj["type"].set<std::string>("registry");
json::value v = json::value(obj);
std::string payload = std::string(v.serialize());
uint32_t len = static_cast<uint32_t>(payload.length());
if ((len + 11U) > DATA_PACKET_LENGTH) {
LogError(LOG_MASTER, "PEER %u (%s) patch status registry payload too large, len = %u", connection->id(), connection->identWithQualifier().c_str(), len);
return false;
}
uint8_t buffer[DATA_PACKET_LENGTH];
::memset(buffer, 0x00U, DATA_PACKET_LENGTH);
::memcpy(buffer + 11U, payload.c_str(), len);
if (m_debug) {
LogDebug(LOG_MASTER, "PEER %u (%s) sending patch status registry, len = %u", connection->id(), connection->identWithQualifier().c_str(), len);
}
return m_host->m_mdNetwork->writePeerMetadata(connection, m_peerId,
{ NET_FUNC::TRANSFER, NET_SUBFUNC::TRANSFER_SUBFUNC_PATCH_STATUS }, buffer, len + 11U, RTP_END_OF_CALL_SEQ, createStreamId());
}
/* Helper to serialize and queue a patch status replication payload. */
bool TrafficNetwork::writePatchStatusReplicationPayload(FNEPeerConnection* connection, json::object obj)
{
if (connection == nullptr)
return false;
if (!m_patchStatusEnabled)
return false;
if (m_host->m_mdNetwork == nullptr)
return false;
if (!connection->connected())
return false;
if (connection->peerClass() != PEER_CONN_CLASS_NEIGHBOR || !connection->isReplica())
return false;
obj["type"].set<std::string>("publish");
json::value v = json::value(obj);
std::string json = std::string(v.serialize());
size_t len = json.length() + 9U;
DECLARE_CHAR_ARRAY(buffer, len);
::memcpy(buffer + 0U, TAG_PEER_REPLICA, 4U);
::snprintf(buffer + 8U, json.length() + 1U, "%s", json.c_str());
PacketBuffer pkt(true, "Peer Replication, Patch Status");
pkt.encode((uint8_t*)buffer, len);
uint32_t streamId = createStreamId();
LogInfoEx(LOG_REPL, "PEER %u (%s) Peer Replication, Patch Status, blocks %u, streamId = %u", connection->id(),
connection->identWithQualifier().c_str(), pkt.fragments.size(), streamId);
if (pkt.fragments.size() > 0U) {
for (auto frag : pkt.fragments) {
m_host->m_mdNetwork->writePeerMetadata(connection, m_peerId, { NET_FUNC::REPL, NET_SUBFUNC::REPL_PATCH_STATUS },
frag.second->data, FRAG_SIZE, RTP_END_OF_CALL_SEQ, streamId);
Thread::sleep(60U); // pace block transmission
}
}
pkt.clear();
return true;
}
/*
** ACL Message Writing
*/

@ -308,35 +308,6 @@ namespace network
* @return json::object
*/
json::object fneConnObject(uint32_t peerId, FNEPeerConnection* conn);
/**
* @brief Gets the console patch status registry.
* @return PatchStatusRegistry& Patch status registry.
*/
PatchStatusRegistry& patchStatusRegistry() { return m_patchStatusRegistry; }
/**
* @brief Flag indicating whether console patch status handling is enabled.
* @returns bool True, if enabled.
*/
bool patchStatusEnabled() const { return m_patchStatusEnabled; }
/**
* @brief Sends patch status registry state to one console peer.
* @param peerId Destination peer ID.
* @param obj Patch status JSON payload.
* @returns bool True, if the message was queued, otherwise false.
*/
bool writePatchStatusToPeer(uint32_t peerId, json::object obj);
/**
* @brief Broadcasts patch status registry state to connected console peers.
* @param obj Patch status JSON payload.
* @param exceptPeerId Optional peer ID to skip.
*/
void writePatchStatusToConsoles(json::object obj, uint32_t exceptPeerId = 0U);
/**
* @brief Replicates patch status state to neighboring FNE peers.
* @param obj Patch status JSON payload.
* @param exceptPeerId Optional peer ID to skip.
*/
void replicatePatchStatus(json::object obj, uint32_t exceptPeerId = 0U);
/**
* @brief Helper to reset a peer connection.
@ -799,6 +770,40 @@ namespace network
*/
static void taskMetadataUpdate(MetadataUpdateRequest* req);
/*
** Console Patch Registry
*/
/**
* @brief Gets the console patch status registry.
* @return PatchStatusRegistry& Patch status registry.
*/
PatchStatusRegistry& patchStatusRegistry() { return m_patchStatusRegistry; }
/**
* @brief Flag indicating whether console patch status handling is enabled.
* @returns bool True, if enabled.
*/
bool patchStatusEnabled() const { return m_patchStatusEnabled; }
/**
* @brief Sends patch status registry state to one console peer.
* @param peerId Destination peer ID.
* @param obj Patch status JSON payload.
* @returns bool True, if the message was queued, otherwise false.
*/
bool writePatchStatusToPeer(uint32_t peerId, json::object obj);
/**
* @brief Broadcasts patch status registry state to connected console peers.
* @param obj Patch status JSON payload.
* @param exceptPeerId Optional peer ID to skip.
*/
void writePatchStatusToConsoles(json::object obj, uint32_t exceptPeerId = 0U);
/**
* @brief Replicates patch status state to neighboring FNE peers.
* @param obj Patch status JSON payload.
* @param exceptPeerId Optional peer ID to skip.
*/
void replicatePatchStatus(json::object obj, uint32_t exceptPeerId = 0U);
/*
** ACL Message Writing
*/

@ -193,7 +193,7 @@ void MetadataNetwork::PacketHandler::transfer(TrafficNetwork* network, MetadataN
FNEPeerConnection* connection = network->m_peers[pktPeerId];
if (connection != nullptr) {
if (!network->patchStatusEnabled()) {
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STATUS, NET_CONN_NAK_FNE_UNAUTHORIZED);
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STAT, NET_CONN_NAK_FNE_UNAUTHORIZED);
break;
}
@ -201,7 +201,7 @@ void MetadataNetwork::PacketHandler::transfer(TrafficNetwork* network, MetadataN
// Only authenticated console peers may publish or request patch registry state.
if (req->length <= TRANSFER_PCKT_HDR_LEN) {
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STATUS, NET_CONN_NAK_ILLEGAL_PACKET);
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STAT, NET_CONN_NAK_ILLEGAL_PACKET);
break;
}
@ -212,7 +212,7 @@ void MetadataNetwork::PacketHandler::transfer(TrafficNetwork* network, MetadataN
json::value v;
std::string err = json::parse(v, payload);
if (!err.empty() || !v.is<json::object>()) {
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STATUS, NET_CONN_NAK_ILLEGAL_PACKET);
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STAT, NET_CONN_NAK_ILLEGAL_PACKET);
break;
}
@ -238,7 +238,7 @@ void MetadataNetwork::PacketHandler::transfer(TrafficNetwork* network, MetadataN
bool changed = false;
if (!network->patchStatusRegistry().publish(reqObj, response, errorMessage, &changed)) {
LogWarning(LOG_MASTER, "PEER %u (%s) invalid patch status payload, %s", pktPeerId, connection->identWithQualifier().c_str(), errorMessage.c_str());
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STATUS, NET_CONN_NAK_ILLEGAL_PACKET);
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STAT, NET_CONN_NAK_ILLEGAL_PACKET);
break;
}
@ -272,7 +272,7 @@ void MetadataNetwork::PacketHandler::transfer(TrafficNetwork* network, MetadataN
}
}
else {
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STATUS, NET_CONN_NAK_FNE_UNAUTHORIZED);
network->writePeerNAK(pktPeerId, network->createStreamId(), TAG_TRANSFER_PATCH_STAT, NET_CONN_NAK_FNE_UNAUTHORIZED);
}
}
}

Loading…
Cancel
Save

Powered by TurnKey Linux.