diff --git a/applications/nucleo-cyphal/CMakeLists.txt b/applications/nucleo-cyphal/CMakeLists.txt index 7cb1152..390f36e 100644 --- a/applications/nucleo-cyphal/CMakeLists.txt +++ b/applications/nucleo-cyphal/CMakeLists.txt @@ -3,6 +3,7 @@ add_firmware(NAME nucleo-cyphal ${CMAKE_CURRENT_SOURCE_DIR}/source/GlobalContext.cpp ${CMAKE_CURRENT_SOURCE_DIR}/source/CyphalApp.cpp ${CMAKE_CURRENT_SOURCE_DIR}/source/GetInfoScanner.cpp + ${CMAKE_CURRENT_SOURCE_DIR}/source/HyphaUdpDispatcher.cpp INCLUDES ${CMAKE_CURRENT_SOURCE_DIR}/include DEFINES diff --git a/applications/nucleo-cyphal/include/HyphaUdpDispatcher.hpp b/applications/nucleo-cyphal/include/HyphaUdpDispatcher.hpp new file mode 100644 index 0000000..2be3b35 --- /dev/null +++ b/applications/nucleo-cyphal/include/HyphaUdpDispatcher.hpp @@ -0,0 +1,52 @@ +#ifndef APP_HYPHA_UDP_DISPATCHER_HPP +#define APP_HYPHA_UDP_DISPATCHER_HPP + +#include + +#include + +extern "C" { +#include "hypha_ip/hypha_ip.h" +} + +#include "jarnax/services/CyphalUDPDispatcher.hpp" + +namespace nucleo { +namespace cyphal { + +/// Adapts libhypha's single UDP callback to Cyphal multicast endpoint dispatch. +class HyphaUdpDispatcher final : public jarnax::cyphal::udp::Dispatcher { +public: + static constexpr std::size_t MaxEndpoints{9U}; + + HyphaUdpDispatcher(HyphaIpContext_t& context, HyphaIpNetworkInterface_t const& network_interface); + + core::Status Join(jarnax::cyphal::udp::Endpoint const& multicast_endpoint, + jarnax::cyphal::udp::DatagramHandler& handler) override; + core::Status Leave(jarnax::cyphal::udp::Endpoint const& multicast_endpoint) override; + core::Status Send(jarnax::cyphal::udp::Endpoint const& destination, + core::Span payload) override; + + /// Delivers a datagram received by CyphalApp's libhypha callback to its registered handler. + void Dispatch(HyphaIpMetaData_t const& metadata, HyphaIpSpan_t datagram); + +private: + struct Registration final { + bool used{false}; + jarnax::cyphal::udp::Endpoint endpoint{}; + jarnax::cyphal::udp::DatagramHandler* handler{nullptr}; + }; + + static HyphaIpIPv4Address_t ToHyphaAddress(std::uint32_t address); + static std::uint32_t ToEndpointAddress(HyphaIpIPv4Address_t address); + static core::Status ToStatus(HyphaIpStatus_e status); + + HyphaIpContext_t& context_; + HyphaIpNetworkInterface_t const& network_interface_; + core::Array registrations_{}; +}; + +} // namespace cyphal +} // namespace nucleo + +#endif // APP_HYPHA_UDP_DISPATCHER_HPP \ No newline at end of file diff --git a/applications/nucleo-cyphal/source/HyphaUdpDispatcher.cpp b/applications/nucleo-cyphal/source/HyphaUdpDispatcher.cpp new file mode 100644 index 0000000..a3e01de --- /dev/null +++ b/applications/nucleo-cyphal/source/HyphaUdpDispatcher.cpp @@ -0,0 +1,96 @@ +#include "HyphaUdpDispatcher.hpp" + +namespace nucleo { +namespace cyphal { + +HyphaUdpDispatcher::HyphaUdpDispatcher(HyphaIpContext_t& context, HyphaIpNetworkInterface_t const& network_interface) + : context_{context} + , network_interface_{network_interface} {} + +HyphaIpIPv4Address_t HyphaUdpDispatcher::ToHyphaAddress(std::uint32_t address) { + return HyphaIpIPv4Address_t{ + static_cast(address >> 24U), static_cast(address >> 16U), + static_cast(address >> 8U), static_cast(address)}; +} + +std::uint32_t HyphaUdpDispatcher::ToEndpointAddress(HyphaIpIPv4Address_t address) { + return (static_cast(address.a) << 24U) | (static_cast(address.b) << 16U) | + (static_cast(address.c) << 8U) | static_cast(address.d); +} + +core::Status HyphaUdpDispatcher::ToStatus(HyphaIpStatus_e status) { + return HyphaIpIsSuccess(status) ? core::Status{} : core::Status{core::Result::Failure, core::Cause::Unknown}; +} + +core::Status HyphaUdpDispatcher::Join(jarnax::cyphal::udp::Endpoint const& multicast_endpoint, + jarnax::cyphal::udp::DatagramHandler& handler) { + for (auto const& registration : registrations_) { + if (registration.used and (registration.endpoint == multicast_endpoint)) { + return core::Status{core::Result::NotExpected, core::Cause::State}; + } + } + Registration* slot = nullptr; + for (auto& registration : registrations_) { + if (not registration.used) { + slot = ®istration; + break; + } + } + if (slot == nullptr) { + return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; + } + core::Status const status = ToStatus( + HyphaIpPrepareUdpReceive(context_, ToHyphaAddress(multicast_endpoint.ip_address), multicast_endpoint.udp_port)); + if (status.IsSuccess()) { + slot->used = true; + slot->endpoint = multicast_endpoint; + slot->handler = &handler; + } + return status; +} + +core::Status HyphaUdpDispatcher::Leave(jarnax::cyphal::udp::Endpoint const& multicast_endpoint) { + for (auto& registration : registrations_) { + if (registration.used and (registration.endpoint == multicast_endpoint)) { + // libhypha has no inverse to HyphaIpPrepareUdpReceive(); prevent local delivery instead. + registration.used = false; + registration.handler = nullptr; + return core::Status{}; + } + } + return core::Status{core::Result::NotExpected, core::Cause::State}; +} + +core::Status HyphaUdpDispatcher::Send(jarnax::cyphal::udp::Endpoint const& destination, + core::Span payload) { + HyphaIpIPv4Address_t const destination_address = ToHyphaAddress(destination.ip_address); + core::Status const prepared = ToStatus(HyphaIpPrepareUdpTransmit(context_, destination_address, destination.udp_port)); + if (not prepared.IsSuccess()) { + return prepared; + } + HyphaIpMetaData_t metadata{}; + metadata.source_address = network_interface_.address; + metadata.destination_address = destination_address; + metadata.source_port = destination.udp_port; + metadata.destination_port = destination.udp_port; + metadata.timestamp = 0U; + HyphaIpSpan_t const datagram{ + const_cast(payload.data()), static_cast(payload.count() & 0x0FFFFFFFU), + HyphaIpSpanTypeUint8_t}; + return ToStatus(HyphaIpTransmitUdpDatagram(context_, &metadata, datagram)); +} + +void HyphaUdpDispatcher::Dispatch(HyphaIpMetaData_t const& metadata, HyphaIpSpan_t datagram) { + jarnax::cyphal::udp::Endpoint const endpoint{ + ToEndpointAddress(metadata.destination_address), metadata.destination_port}; + for (auto const& registration : registrations_) { + if (registration.used and (registration.endpoint == endpoint) and (registration.handler != nullptr)) { + registration.handler->OnDatagramReceived( + endpoint, static_cast(datagram.pointer), HyphaIpSpanSize(datagram)); + return; + } + } +} + +} // namespace cyphal +} // namespace nucleo \ No newline at end of file diff --git a/modules/jarnax/source/include/jarnax/services/CyphalUDPSocket.hpp b/modules/jarnax/source/include/jarnax/services/CyphalUDPDispatcher.hpp similarity index 77% rename from modules/jarnax/source/include/jarnax/services/CyphalUDPSocket.hpp rename to modules/jarnax/source/include/jarnax/services/CyphalUDPDispatcher.hpp index 159a0d3..f436d56 100644 --- a/modules/jarnax/source/include/jarnax/services/CyphalUDPSocket.hpp +++ b/modules/jarnax/source/include/jarnax/services/CyphalUDPDispatcher.hpp @@ -1,5 +1,5 @@ -#ifndef JARNAX_SERVICES_CYPHAL_UDP_SOCKET_HPP -#define JARNAX_SERVICES_CYPHAL_UDP_SOCKET_HPP +#ifndef JARNAX_SERVICES_CYPHAL_UDP_DISPATCHER_HPP +#define JARNAX_SERVICES_CYPHAL_UDP_DISPATCHER_HPP #include #include @@ -26,7 +26,7 @@ struct Endpoint final { } }; -/// The handler of datagrams received by a Socket. +/// The handler of datagrams received by a Dispatcher. /// Implementations are notified for every datagram the underlying UDP/IP stack delivers. class DatagramHandler { public: @@ -40,27 +40,26 @@ class DatagramHandler { ~DatagramHandler() = default; }; -/// The minimal abstraction over the UDP/IP stack required by Cyphal/UDP. +/// The narrow UDP/IP dispatch boundary required by Cyphal/UDP. /// The production implementation wraps hypha; unit tests inject a mock. /// All methods are non-blocking and return immediately. -class Socket { +class Dispatcher { public: - /// Joins the multicast group at the endpoint and begins delivering received - /// datagrams to the handler (e.g. IGMP join plus socket bind). + /// Begins delivering multicast datagrams for the endpoint to the handler. virtual core::Status Join(Endpoint const& multicast_endpoint, DatagramHandler& handler) = 0; - /// Leaves the multicast group at the endpoint; no more datagrams are delivered. + /// Stops delivering multicast datagrams for the endpoint. virtual core::Status Leave(Endpoint const& multicast_endpoint) = 0; /// Transmits one datagram to the given multicast endpoint. virtual core::Status Send(Endpoint const& destination, core::Span payload) = 0; protected: - ~Socket() = default; + ~Dispatcher() = default; }; } // namespace udp } // namespace cyphal } // namespace jarnax -#endif // JARNAX_SERVICES_CYPHAL_UDP_SOCKET_HPP +#endif // JARNAX_SERVICES_CYPHAL_UDP_DISPATCHER_HPP \ No newline at end of file diff --git a/modules/jarnax/source/include/jarnax/services/CyphalUDPInterface.hpp b/modules/jarnax/source/include/jarnax/services/CyphalUDPInterface.hpp index c1b88ff..8788c6b 100644 --- a/modules/jarnax/source/include/jarnax/services/CyphalUDPInterface.hpp +++ b/modules/jarnax/source/include/jarnax/services/CyphalUDPInterface.hpp @@ -14,7 +14,7 @@ extern "C" { #include "jarnax/Loopable.hpp" #include "jarnax/cyphal/Interface.hpp" #include "jarnax/cyphal/O1HeapPool.hpp" -#include "jarnax/services/CyphalUDPSocket.hpp" +#include "jarnax/services/CyphalUDPDispatcher.hpp" namespace jarnax { namespace cyphal { @@ -33,7 +33,7 @@ class MicrosecondClock { /// The concrete Cyphal/UDP transport implementing cyphal::Interface on top of libudpard. /// Applications never touch udpard types directly; they use Listen/Send/RegisterListener /// and drive the object from the SuperLoop via Loopable::Execute (TX drain). -/// Incoming datagrams are pushed in by a udp::Socket implementation (hypha on target). +/// Incoming datagrams are pushed in by a udp::Dispatcher implementation (hypha on target). class CyphalUDPInterface final : public Interface, public Loopable, public udp::DatagramHandler { public: /// The maximum number of subjects which can be listened to concurrently. @@ -50,9 +50,9 @@ class CyphalUDPInterface final : public Interface, public Loopable, public udp:: /// Constructs the interface. /// @param heap The O1Heap backed pool shared by all udpard allocations. /// @param node_id The local Cyphal node-ID; anonymous (0xFFFF) cannot use services. - /// @param socket The UDP/IP stack abstraction; its lifetime must exceed ours. + /// @param dispatcher The UDP/IP dispatcher; its lifetime must exceed ours. /// @param ticker The time source used for transfer deadlines and timestamps. - CyphalUDPInterface(O1HeapPool& heap, udp::NodeId node_id, udp::Socket& socket, MicrosecondClock& clock); + CyphalUDPInterface(O1HeapPool& heap, udp::NodeId node_id, udp::Dispatcher& dispatcher, MicrosecondClock& clock); virtual ~CyphalUDPInterface() override; //+=== LOOPABLE INTERFACE === @@ -118,7 +118,7 @@ class CyphalUDPInterface final : public Interface, public Loopable, public udp:: void DeliverTransfer(UdpardRxTransfer const& transfer, PortId port_id); O1HeapPool& heap_; - udp::Socket& socket_; + udp::Dispatcher& dispatcher_; MicrosecondClock& clock_; udp::NodeId local_node_id_; bool initialized_; @@ -126,7 +126,7 @@ class CyphalUDPInterface final : public Interface, public Loopable, public udp:: UdpardMemoryResource tx_memory_; UdpardRxMemoryResources rx_memory_; UdpardTx tx_; - UdpardRxRPCDispatcher dispatcher_; + UdpardRxRPCDispatcher rpc_dispatcher_; udp::Endpoint service_endpoint_{}; bool service_group_joined_{false}; diff --git a/modules/jarnax/source/services/CyphalUDPInterface.cpp b/modules/jarnax/source/services/CyphalUDPInterface.cpp index d8d5b58..3259f4c 100644 --- a/modules/jarnax/source/services/CyphalUDPInterface.cpp +++ b/modules/jarnax/source/services/CyphalUDPInterface.cpp @@ -15,8 +15,8 @@ constexpr bool SamePort(PortId const& lhs, PortId const& rhs) { return false; } bool const same = (lhs.type == PortId::Type::Subject) // - ? (lhs.subject.value == rhs.subject.value) - : ((lhs.style == rhs.style) and (lhs.service.value == rhs.service.value)); + ? (lhs.subject.value == rhs.subject.value) + : ((lhs.style == rhs.style) and (lhs.service.value == rhs.service.value)); return same; } @@ -31,24 +31,29 @@ core::Status CyphalUDPInterface::MapResult(std::int32_t result) const { return core::Status{}; // Success } switch (-result) { - case UDPARD_ERROR_ARGUMENT: return core::Status{core::Result::InvalidValue, core::Cause::Parameter}; - case UDPARD_ERROR_MEMORY: return core::Status{core::Result::NotEnough, core::Cause::Resource}; - case UDPARD_ERROR_CAPACITY: return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; - case UDPARD_ERROR_ANONYMOUS: return core::Status{core::Result::NotConfigured, core::Cause::Configuration}; - default: return core::Status{core::Result::Failure, core::Cause::Unknown}; + case UDPARD_ERROR_ARGUMENT: + return core::Status{core::Result::InvalidValue, core::Cause::Parameter}; + case UDPARD_ERROR_MEMORY: + return core::Status{core::Result::NotEnough, core::Cause::Resource}; + case UDPARD_ERROR_CAPACITY: + return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; + case UDPARD_ERROR_ANONYMOUS: + return core::Status{core::Result::NotConfigured, core::Cause::Configuration}; + default: + return core::Status{core::Result::Failure, core::Cause::Unknown}; } } -CyphalUDPInterface::CyphalUDPInterface(O1HeapPool& heap, udp::NodeId node_id, udp::Socket& socket, MicrosecondClock& clock) +CyphalUDPInterface::CyphalUDPInterface(O1HeapPool& heap, udp::NodeId node_id, udp::Dispatcher& dispatcher, MicrosecondClock& clock) : heap_{heap} - , socket_{socket} + , dispatcher_{dispatcher} , clock_{clock} , local_node_id_{node_id} , initialized_{false} , tx_memory_{} , rx_memory_{} , tx_{} - , dispatcher_{} + , rpc_dispatcher_{} , service_endpoint_{} , service_group_joined_{false} , listener_{nullptr} @@ -62,13 +67,13 @@ CyphalUDPInterface::CyphalUDPInterface(O1HeapPool& heap, udp::NodeId node_id, ud } tx_.mtu = UDPARD_MTU_DEFAULT; - std::int_fast8_t const dispatcher_init = udpardRxRPCDispatcherInit(&dispatcher_, rx_memory_); + std::int_fast8_t const dispatcher_init = udpardRxRPCDispatcherInit(&rpc_dispatcher_, rx_memory_); if (dispatcher_init < 0) { return; } UdpardUDPIPEndpoint endpoint{}; - std::int_fast8_t const dispatcher_start = udpardRxRPCDispatcherStart(&dispatcher_, local_node_id_, &endpoint); + std::int_fast8_t const dispatcher_start = udpardRxRPCDispatcherStart(&rpc_dispatcher_, local_node_id_, &endpoint); if (dispatcher_start < 0) { return; } @@ -87,7 +92,7 @@ CyphalUDPInterface::~CyphalUDPInterface() { } for (auto& port : service_ports_) { if (port.used) { - (void)udpardRxRPCDispatcherCancel(&dispatcher_, port.service_id.value, port.is_request); + (void)udpardRxRPCDispatcherCancel(&rpc_dispatcher_, port.service_id.value, port.is_request); port.used = false; } } @@ -104,10 +109,8 @@ bool CyphalUDPInterface::Execute() { // Drain the prioritized TX queue, highest priority first. while (UdpardTxItem const* item = udpardTxPeek(&tx_)) { udp::Endpoint const destination{item->destination.ip_address, item->destination.udp_port}; - core::Span const payload{ - static_cast(item->datagram_payload.data), item->datagram_payload.size - }; - core::Status const status = socket_.Send(destination, payload); + core::Span const payload{static_cast(item->datagram_payload.data), item->datagram_payload.size}; + core::Status const status = dispatcher_.Send(destination, payload); UdpardTxItem* const taken = udpardTxPop(&tx_, item); udpardTxFree(tx_memory_, taken); if (not status.IsSuccess()) { @@ -173,13 +176,12 @@ core::Status CyphalUDPInterface::ListenSubject(PortId port_id) { if (slot == nullptr) { return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; } - std::int_fast8_t const result = - udpardRxSubscriptionInit(&slot->sub, static_cast(subject.value), MaxExtent, rx_memory_); + std::int_fast8_t const result = udpardRxSubscriptionInit(&slot->sub, static_cast(subject.value), MaxExtent, rx_memory_); if (result < 0) { return MapResult(result); } udp::Endpoint const group{slot->sub.udp_ip_endpoint.ip_address, slot->sub.udp_ip_endpoint.udp_port}; - core::Status const join = socket_.Join(group, *this); + core::Status const join = dispatcher_.Join(group, *this); if (not join.IsSuccess()) { udpardRxSubscriptionFree(&slot->sub); return join; @@ -200,14 +202,14 @@ core::Status CyphalUDPInterface::ListenService(PortId port_id) { return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; } std::int_fast8_t const listen = - udpardRxRPCDispatcherListen(&dispatcher_, &slot->port, static_cast(service.value), is_request, MaxExtent); + udpardRxRPCDispatcherListen(&rpc_dispatcher_, &slot->port, static_cast(service.value), is_request, MaxExtent); if (listen < 0) { return MapResult(listen); } if (not service_group_joined_) { - core::Status const join = socket_.Join(service_endpoint_, *this); + core::Status const join = dispatcher_.Join(service_endpoint_, *this); if (not join.IsSuccess()) { - (void)udpardRxRPCDispatcherCancel(&dispatcher_, static_cast(service.value), is_request); + (void)udpardRxRPCDispatcherCancel(&rpc_dispatcher_, static_cast(service.value), is_request); return join; } service_group_joined_ = true; @@ -240,7 +242,7 @@ core::Status CyphalUDPInterface::RemoveSubject(PortId port_id) { udp::Endpoint const group{slot->sub.udp_ip_endpoint.ip_address, slot->sub.udp_ip_endpoint.udp_port}; udpardRxSubscriptionFree(&slot->sub); slot->used = false; - return socket_.Leave(group); + return dispatcher_.Leave(group); } core::Status CyphalUDPInterface::RemoveService(PortId port_id) { @@ -250,8 +252,7 @@ core::Status CyphalUDPInterface::RemoveService(PortId port_id) { if (slot == nullptr) { return core::Status{core::Result::NotExpected, core::Cause::State}; // not listening } - std::int_fast8_t const cancel = - udpardRxRPCDispatcherCancel(&dispatcher_, static_cast(service.value), is_request); + std::int_fast8_t const cancel = udpardRxRPCDispatcherCancel(&rpc_dispatcher_, static_cast(service.value), is_request); slot->used = false; if (cancel < 0) { return MapResult(cancel); @@ -263,7 +264,7 @@ core::Status CyphalUDPInterface::RemoveService(PortId port_id) { } if (not any_left and service_group_joined_) { service_group_joined_ = false; - return socket_.Leave(service_endpoint_); + return dispatcher_.Leave(service_endpoint_); } return core::Status{}; } @@ -304,8 +305,7 @@ CyphalUDPInterface::TransferIdCounter* CyphalUDPInterface::FindTransferCounter(P } continue; } - bool const same_value = (counter.value == - ((port_id.type == PortId::Type::Subject) ? port_id.subject.value : port_id.service.value)); + bool const same_value = (counter.value == ((port_id.type == PortId::Type::Subject) ? port_id.subject.value : port_id.service.value)); bool const same_kind = (counter.type == port_id.type) and (counter.style == port_id.style); bool const same_peer = (counter.type == PortId::Type::Subject) or (counter.peer == peer); if (same_value and same_kind and same_peer) { @@ -316,16 +316,14 @@ CyphalUDPInterface::TransferIdCounter* CyphalUDPInterface::FindTransferCounter(P free_slot->used = true; free_slot->type = port_id.type; free_slot->style = port_id.style; - free_slot->value = - (port_id.type == PortId::Type::Subject) ? port_id.subject.value : port_id.service.value; + free_slot->value = (port_id.type == PortId::Type::Subject) ? port_id.subject.value : port_id.service.value; free_slot->peer = peer; free_slot->next = 0U; } return free_slot; } -CyphalUDPInterface::PendingRequest* CyphalUDPInterface::RememberRequest( - ServiceId service_id, udp::NodeId client, UdpardTransferID transfer_id) { +CyphalUDPInterface::PendingRequest* CyphalUDPInterface::RememberRequest(ServiceId service_id, udp::NodeId client, UdpardTransferID transfer_id) { PendingRequest* slot = const_cast(FindPendingRequest(service_id, client)); if (slot == nullptr) { for (auto& pending : pending_requests_) { @@ -346,8 +344,7 @@ CyphalUDPInterface::PendingRequest* CyphalUDPInterface::RememberRequest( return slot; } -CyphalUDPInterface::PendingRequest const* CyphalUDPInterface::FindPendingRequest( - ServiceId service_id, udp::NodeId client) const { +CyphalUDPInterface::PendingRequest const* CyphalUDPInterface::FindPendingRequest(ServiceId service_id, udp::NodeId client) const { for (auto const& pending : pending_requests_) { if (pending.used and (pending.service_id.value == service_id.value) and (pending.client == client)) { return &pending; @@ -370,14 +367,12 @@ core::Status CyphalUDPInterface::Send(Metadata& metadata, SerializedMessage msg) return core::Status{core::Result::ExceededLimit, core::Cause::Resource}; } UdpardTransferID const transfer_id = counter->next; - result = udpardTxPublish( - &tx_, deadline, DefaultPriority, static_cast(metadata.port_id.subject.value), transfer_id, payload, - this); + result = + udpardTxPublish(&tx_, deadline, DefaultPriority, static_cast(metadata.port_id.subject.value), transfer_id, payload, this); if (result > 0) { counter->next = transfer_id + 1U; // increment only on success per the library contract } - } else if ( - (metadata.port_id.type == PortId::Type::Service) and (metadata.port_id.style == PortId::Style::Request)) { + } else if ((metadata.port_id.type == PortId::Type::Service) and (metadata.port_id.style == PortId::Style::Request)) { ServiceId const service = metadata.port_id.service; TransferIdCounter* const counter = FindTransferCounter(metadata.port_id, metadata.recipient); if (counter == nullptr) { @@ -385,21 +380,34 @@ core::Status CyphalUDPInterface::Send(Metadata& metadata, SerializedMessage msg) } UdpardTransferID const transfer_id = counter->next; result = udpardTxRequest( - &tx_, deadline, DefaultPriority, static_cast(service.value), - static_cast(metadata.recipient), transfer_id, payload, this); + &tx_, + deadline, + DefaultPriority, + static_cast(service.value), + static_cast(metadata.recipient), + transfer_id, + payload, + this + ); if (result > 0) { counter->next = transfer_id + 1U; } - } else if ( - (metadata.port_id.type == PortId::Type::Service) and (metadata.port_id.style == PortId::Style::Response)) { + } else if ((metadata.port_id.type == PortId::Type::Service) and (metadata.port_id.style == PortId::Style::Response)) { ServiceId const service = metadata.port_id.service; PendingRequest const* const pending = FindPendingRequest(service, metadata.recipient); if (pending == nullptr) { return core::Status{core::Result::NotExpected, core::Cause::State}; // no request to respond to } result = udpardTxRespond( - &tx_, deadline, DefaultPriority, static_cast(service.value), - static_cast(metadata.recipient), pending->transfer_id, payload, this); + &tx_, + deadline, + DefaultPriority, + static_cast(service.value), + static_cast(metadata.recipient), + pending->transfer_id, + payload, + this + ); } else { return core::Status{core::Result::InvalidValue, core::Cause::Parameter}; } @@ -422,20 +430,15 @@ void CyphalUDPInterface::DeliverTransfer(UdpardRxTransfer const& transfer, PortI return; } // Empty transfers (e.g. GetInfo requests) are valid and delivered as zero-length messages. - size_t const expected = - (transfer.payload_size < rx_scratch_.size()) ? transfer.payload_size : rx_scratch_.size(); + size_t const expected = (transfer.payload_size < rx_scratch_.size()) ? transfer.payload_size : rx_scratch_.size(); if (expected > 0U) { size_t const gathered = udpardGather(transfer.payload, expected, rx_scratch_.data()); if (gathered == 0U) { return; } } - NodeId const recipient = - (port_id.type == PortId::Type::Service) ? local_node_id_ : udp::anonymous; - Metadata const metadata{ - static_cast(transfer.source_node_id), recipient, port_id, - core::units::MicroSeconds{transfer.timestamp_usec} - }; + NodeId const recipient = (port_id.type == PortId::Type::Service) ? local_node_id_ : udp::anonymous; + Metadata const metadata{static_cast(transfer.source_node_id), recipient, port_id, core::units::MicroSeconds{transfer.timestamp_usec}}; SerializedMessage const msg{rx_scratch_.data(), expected}; listener_->OnReceive(metadata, msg); statistics_.transfer.num_received++; @@ -458,16 +461,12 @@ void CyphalUDPInterface::OnDatagramReceived(udp::Endpoint const& destination, st if (destination == service_endpoint_) { UdpardRxRPCTransfer transfer{}; - std::int_fast8_t const result = - udpardRxRPCDispatcherReceive(&dispatcher_, NowUs(), payload, 0U, nullptr, &transfer); + std::int_fast8_t const result = udpardRxRPCDispatcherReceive(&rpc_dispatcher_, NowUs(), payload, 0U, nullptr, &transfer); if (result > 0) { ServiceId const service{static_cast(transfer.service_id)}; if (transfer.is_request) { - (void)RememberRequest( - service, static_cast(transfer.base.source_node_id), transfer.base.transfer_id); - DeliverTransfer( - transfer.base, - PortId{service, PortId::Style::Request}); + (void)RememberRequest(service, static_cast(transfer.base.source_node_id), transfer.base.transfer_id); + DeliverTransfer(transfer.base, PortId{service, PortId::Style::Request}); } else { DeliverTransfer(transfer.base, PortId{service, PortId::Style::Response}); } @@ -476,8 +475,8 @@ void CyphalUDPInterface::OnDatagramReceived(udp::Endpoint const& destination, st } else { Subscription* slot = nullptr; for (auto& sub : subscriptions_) { - if (sub.used // - and (sub.sub.udp_ip_endpoint.ip_address == destination.ip_address) // + if (sub.used // + and (sub.sub.udp_ip_endpoint.ip_address == destination.ip_address) // and (sub.sub.udp_ip_endpoint.udp_port == destination.udp_port)) { slot = ⊂ break; @@ -485,8 +484,7 @@ void CyphalUDPInterface::OnDatagramReceived(udp::Endpoint const& destination, st } if (slot != nullptr) { UdpardRxTransfer transfer{}; - std::int_fast8_t const result = - udpardRxSubscriptionReceive(&slot->sub, NowUs(), payload, 0U, &transfer); + std::int_fast8_t const result = udpardRxSubscriptionReceive(&slot->sub, NowUs(), payload, 0U, &transfer); if (result > 0) { DeliverTransfer(transfer, PortId{slot->subject_id}); udpardRxFragmentFree(transfer.payload, rx_memory_.fragment, rx_memory_.payload); diff --git a/modules/jarnax/tests/gtest-cyphal-udpinterface.cpp b/modules/jarnax/tests/gtest-cyphal-udpinterface.cpp index 3ab1f7d..12b8d06 100644 --- a/modules/jarnax/tests/gtest-cyphal-udpinterface.cpp +++ b/modules/jarnax/tests/gtest-cyphal-udpinterface.cpp @@ -11,7 +11,7 @@ extern "C" { #include "core/units/MicroSeconds.hpp" #include "jarnax/cyphal/O1HeapPool.hpp" #include "jarnax/services/CyphalUDPInterface.hpp" -#include "jarnax/services/MockUDPSocket.hpp" +#include "jarnax/services/MockUDPDispatcher.hpp" namespace { @@ -26,7 +26,7 @@ using jarnax::cyphal::SerializedMessage; using jarnax::cyphal::ServiceId; using jarnax::cyphal::SubjectId; using jarnax::cyphal::udp::Endpoint; -using jarnax::cyphal::udp::MockSocket; +using jarnax::cyphal::udp::MockDispatcher; constexpr std::uint16_t LocalNodeId{42U}; constexpr std::uint16_t RemoteNodeId{100U}; @@ -95,8 +95,7 @@ class DatagramProducer final { /// Publishes a message on a subject; returns the datagram to feed into the interface. bool Publish(std::uint16_t subject_id, UdpardTransferID transfer_id, std::vector const& message) { UdpardPayload const payload{message.size(), message.data()}; - std::int32_t const result = - udpardTxPublish(&tx_, 1000000ULL, UdpardPriorityNominal, subject_id, transfer_id, payload, nullptr); + std::int32_t const result = udpardTxPublish(&tx_, 1000000ULL, UdpardPriorityNominal, subject_id, transfer_id, payload, nullptr); return result > 0; } @@ -104,18 +103,17 @@ class DatagramProducer final { bool Request(std::uint16_t service_id, std::uint16_t server_node, UdpardTransferID transfer_id) { UdpardPayload const payload{0U, nullptr}; std::int32_t const result = udpardTxRequest( - &tx_, 1000000ULL, UdpardPriorityNominal, service_id, static_cast(server_node), transfer_id, - payload, nullptr); + &tx_, 1000000ULL, UdpardPriorityNominal, service_id, static_cast(server_node), transfer_id, payload, nullptr + ); return result > 0; } /// Sends an RPC response to the given client node. - bool Respond(std::uint16_t service_id, std::uint16_t client_node, UdpardTransferID transfer_id, - std::vector const& message) { + bool Respond(std::uint16_t service_id, std::uint16_t client_node, UdpardTransferID transfer_id, std::vector const& message) { UdpardPayload const payload{message.size(), message.data()}; std::int32_t const result = udpardTxRespond( - &tx_, 1000000ULL, UdpardPriorityNominal, service_id, static_cast(client_node), transfer_id, - payload, nullptr); + &tx_, 1000000ULL, UdpardPriorityNominal, service_id, static_cast(client_node), transfer_id, payload, nullptr + ); return result > 0; } @@ -143,12 +141,12 @@ class DatagramProducer final { class CyphalUDPInterfaceTest : public Test { protected: void SetUp() override { - interface_ = std::make_unique(heap_, static_cast(LocalNodeId), socket_, clock_); + interface_ = std::make_unique(heap_, static_cast(LocalNodeId), dispatcher_, clock_); ASSERT_TRUE(interface_->IsInitialized()); } jarnax::cyphal::O1HeapPool& heap_{jarnax::cyphal::O1HeapPool::Instance()}; - NiceMock socket_{}; + NiceMock dispatcher_{}; FakeClock clock_; TestListener listener_; std::unique_ptr interface_; @@ -166,20 +164,17 @@ TEST_F(CyphalUDPInterfaceTest, RegisterListenerSucceedsAndReplaces) { TEST_F(CyphalUDPInterfaceTest, ListenSubjectJoinsMulticastGroup) { Endpoint joined{}; - EXPECT_CALL(socket_, Join(_, _)) - .WillOnce(DoAll(SaveArg<0>(&joined), Return(core::Status{}))) - .WillRepeatedly(Return(core::Status{})); + EXPECT_CALL(dispatcher_, Join(_, _)).WillOnce(DoAll(SaveArg<0>(&joined), Return(core::Status{}))).WillRepeatedly(Return(core::Status{})); EXPECT_TRUE(interface_->Listen(PortId{TestSubject}).IsSuccess()); EXPECT_TRUE(interface_->IsListening(PortId{TestSubject})); EXPECT_NE(joined.udp_port, 0U); // derived from the subject by libudpard // Duplicate subscription is rejected - EXPECT_EQ( - interface_->Listen(PortId{TestSubject}).GetResult(), core::Result::NotExpected); + EXPECT_EQ(interface_->Listen(PortId{TestSubject}).GetResult(), core::Result::NotExpected); // Removal leaves the group - EXPECT_CALL(socket_, Leave(_)).WillOnce(Return(core::Status{})); + EXPECT_CALL(dispatcher_, Leave(_)).WillOnce(Return(core::Status{})); EXPECT_TRUE(interface_->Remove(PortId{TestSubject}).IsSuccess()); EXPECT_FALSE(interface_->IsListening(PortId{TestSubject})); @@ -187,17 +182,15 @@ TEST_F(CyphalUDPInterfaceTest, ListenSubjectJoinsMulticastGroup) { EXPECT_EQ(interface_->Remove(PortId{TestSubject}).GetResult(), core::Result::NotExpected); // Verify expectations at this checkpoint - Mock::VerifyAndClearExpectations(&socket_); + Mock::VerifyAndClearExpectations(&dispatcher_); } TEST_F(CyphalUDPInterfaceTest, ListenServicePortsJoinServiceGroupOnce) { Endpoint service_endpoint{}; - EXPECT_CALL(socket_, Join(_, _)) - .Times(1) - .WillOnce([&service_endpoint](Endpoint const& endpoint, jarnax::cyphal::udp::DatagramHandler&) { - service_endpoint = endpoint; - return core::Status{}; - }); + EXPECT_CALL(dispatcher_, Join(_, _)).Times(1).WillOnce([&service_endpoint](Endpoint const& endpoint, jarnax::cyphal::udp::DatagramHandler&) { + service_endpoint = endpoint; + return core::Status{}; + }); auto const request = PortId{GetInfoServiceId, PortId::Style::Request}; auto const response = PortId{GetInfoServiceId, PortId::Style::Response}; @@ -212,52 +205,47 @@ TEST_F(CyphalUDPInterfaceTest, ListenServicePortsJoinServiceGroupOnce) { EXPECT_EQ(interface_->Listen(request).GetResult(), core::Result::NotExpected); EXPECT_EQ(interface_->Listen(response).GetResult(), core::Result::NotExpected); - EXPECT_CALL(socket_, Leave(_)).WillOnce(Return(core::Status{})); - EXPECT_TRUE(interface_->Remove(response).IsSuccess()); // group still needed by request port - EXPECT_TRUE(interface_->Remove(request).IsSuccess()); // last port leaves the group + EXPECT_CALL(dispatcher_, Leave(_)).WillOnce(Return(core::Status{})); + EXPECT_TRUE(interface_->Remove(response).IsSuccess()); // group still needed by request port + EXPECT_TRUE(interface_->Remove(request).IsSuccess()); // last port leaves the group EXPECT_FALSE(interface_->IsListening(request)); EXPECT_FALSE(interface_->IsListening(response)); - Mock::VerifyAndClearExpectations(&socket_); + Mock::VerifyAndClearExpectations(&dispatcher_); } TEST_F(CyphalUDPInterfaceTest, ListenRejectsInvalidPorts) { // A service port explicitly styled as Neither is invalid - EXPECT_EQ( - interface_->Listen(PortId{GetInfoServiceId, PortId::Style::Neither}).GetResult(), core::Result::InvalidValue); + EXPECT_EQ(interface_->Listen(PortId{GetInfoServiceId, PortId::Style::Neither}).GetResult(), core::Result::InvalidValue); } TEST_F(CyphalUDPInterfaceTest, ListenFailsWhenSubscriptionsExhausted) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); for (std::size_t i = 0U; i < CyphalUDPInterface::MaxSubscriptions; ++i) { EXPECT_TRUE(interface_->Listen(PortId{SubjectId{static_cast(100U + i)}}).IsSuccess()); } - EXPECT_EQ( - interface_->Listen(PortId{SubjectId{static_cast(200U)}}).GetResult(), - core::Result::ExceededLimit); + EXPECT_EQ(interface_->Listen(PortId{SubjectId{static_cast(200U)}}).GetResult(), core::Result::ExceededLimit); } TEST_F(CyphalUDPInterfaceTest, SendPublishesToSubjectMulticastGroup) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); EXPECT_TRUE(interface_->Listen(PortId{TestSubject}).IsSuccess()); std::vector message{1U, 2U, 3U, 4U}; - Metadata metadata{LocalNodeId, jarnax::cyphal::udp::anonymous, PortId{TestSubject}, - core::units::MicroSeconds{0ULL}}; + Metadata metadata{LocalNodeId, jarnax::cyphal::udp::anonymous, PortId{TestSubject}, core::units::MicroSeconds{0ULL}}; EXPECT_TRUE(interface_->Send(metadata, SerializedMessage{message.data(), message.size()}).IsSuccess()); // Drain: exactly one datagram goes out on the subject multicast group Endpoint sent_to{}; std::size_t payload_size = 0U; - EXPECT_CALL(socket_, Send(_, _)) - .WillOnce([&sent_to, &payload_size](Endpoint const& destination, core::Span payload) { - sent_to = destination; - payload_size = payload.count(); - return core::Status{}; - }); + EXPECT_CALL(dispatcher_, Send(_, _)).WillOnce([&sent_to, &payload_size](Endpoint const& destination, core::Span payload) { + sent_to = destination; + payload_size = payload.count(); + return core::Status{}; + }); EXPECT_TRUE(interface_->Execute()); - EXPECT_EQ(sent_to.udp_port, 9382U); // the Cyphal/UDP well-known port + EXPECT_EQ(sent_to.udp_port, 9382U); // the Cyphal/UDP well-known port EXPECT_GT(payload_size, message.size()); // header plus CRC are added jarnax::cyphal::TransportStatistics statistics{}; @@ -266,19 +254,18 @@ TEST_F(CyphalUDPInterfaceTest, SendPublishesToSubjectMulticastGroup) { EXPECT_EQ(statistics.transfer.num_emitted, 1U); EXPECT_EQ(statistics.network_interfaces[0U].num_emitted, 1U); - Mock::VerifyAndClearExpectations(&socket_); + Mock::VerifyAndClearExpectations(&dispatcher_); } TEST_F(CyphalUDPInterfaceTest, ExecuteReportsTransmitErrors) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); EXPECT_TRUE(interface_->Listen(PortId{TestSubject}).IsSuccess()); std::vector message{1U}; - Metadata metadata{LocalNodeId, jarnax::cyphal::udp::anonymous, PortId{TestSubject}, - core::units::MicroSeconds{0ULL}}; + Metadata metadata{LocalNodeId, jarnax::cyphal::udp::anonymous, PortId{TestSubject}, core::units::MicroSeconds{0ULL}}; EXPECT_TRUE(interface_->Send(metadata, SerializedMessage{message.data(), message.size()}).IsSuccess()); - EXPECT_CALL(socket_, Send(_, _)).WillOnce(Return(core::Status{core::Result::Failure, core::Cause::Peripheral})); + EXPECT_CALL(dispatcher_, Send(_, _)).WillOnce(Return(core::Status{core::Result::Failure, core::Cause::Peripheral})); EXPECT_TRUE(interface_->Execute()); jarnax::cyphal::TransportStatistics statistics{}; @@ -286,11 +273,11 @@ TEST_F(CyphalUDPInterfaceTest, ExecuteReportsTransmitErrors) { EXPECT_EQ(statistics.transfer.num_emitted, 0U); EXPECT_EQ(statistics.transfer.num_errored, 1U); - Mock::VerifyAndClearExpectations(&socket_); + Mock::VerifyAndClearExpectations(&dispatcher_); } TEST_F(CyphalUDPInterfaceTest, ReceiveDeliversSubjectTransferToListener) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); EXPECT_TRUE(interface_->Listen(PortId{TestSubject}).IsSuccess()); EXPECT_TRUE(interface_->RegisterListener(LocalNodeId, listener_).IsSuccess()); @@ -319,7 +306,7 @@ TEST_F(CyphalUDPInterfaceTest, ReceiveDeliversSubjectTransferToListener) { } TEST_F(CyphalUDPInterfaceTest, ReceiveIgnoresUnknownGroupsAndBadDatagrams) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); EXPECT_TRUE(interface_->Listen(PortId{TestSubject}).IsSuccess()); EXPECT_TRUE(interface_->RegisterListener(LocalNodeId, listener_).IsSuccess()); @@ -348,7 +335,7 @@ TEST_F(CyphalUDPInterfaceTest, ReceiveIgnoresUnknownGroupsAndBadDatagrams) { } TEST_F(CyphalUDPInterfaceTest, ServiceRequestResponseRoundTrip) { - ON_CALL(socket_, Join(_, _)).WillByDefault(Return(core::Status{})); + ON_CALL(dispatcher_, Join(_, _)).WillByDefault(Return(core::Status{})); auto const request_port = PortId{GetInfoServiceId, PortId::Style::Request}; auto const response_port = PortId{GetInfoServiceId, PortId::Style::Response}; EXPECT_TRUE(interface_->Listen(request_port).IsSuccess()); @@ -373,16 +360,14 @@ TEST_F(CyphalUDPInterfaceTest, ServiceRequestResponseRoundTrip) { // We respond via the Interface API; the recorded request transfer-ID is echoed. std::vector response_message{5U, 6U}; Metadata metadata{LocalNodeId, RemoteNodeId, response_port, core::units::MicroSeconds{0ULL}}; - EXPECT_TRUE(interface_->Send(metadata, SerializedMessage{response_message.data(), response_message.size()}) - .IsSuccess()); + EXPECT_TRUE(interface_->Send(metadata, SerializedMessage{response_message.data(), response_message.size()}).IsSuccess()); // The response goes to the client's RPC multicast group when drained. Endpoint sent_to{}; - EXPECT_CALL(socket_, Send(_, _)) - .WillOnce([&](Endpoint const& destination, core::Span) { - sent_to = destination; - return core::Status{}; - }); + EXPECT_CALL(dispatcher_, Send(_, _)).WillOnce([&](Endpoint const& destination, core::Span) { + sent_to = destination; + return core::Status{}; + }); EXPECT_TRUE(interface_->Execute()); // The response is addressed to the client's own RPC multicast group (node 100). @@ -406,15 +391,14 @@ TEST_F(CyphalUDPInterfaceTest, ServiceRequestResponseRoundTrip) { EXPECT_TRUE(interface_->GetStatistics(statistics).IsSuccess()); EXPECT_EQ(statistics.transfer.num_received, 2U); - Mock::VerifyAndClearExpectations(&socket_); + Mock::VerifyAndClearExpectations(&dispatcher_); } TEST_F(CyphalUDPInterfaceTest, SendResponseWithoutRequestFails) { auto const response_port = PortId{GetInfoServiceId, PortId::Style::Response}; Metadata metadata{LocalNodeId, RemoteNodeId, response_port, core::units::MicroSeconds{0ULL}}; std::uint8_t byte{0U}; - EXPECT_EQ( - interface_->Send(metadata, SerializedMessage{&byte, 1U}).GetResult(), core::Result::NotExpected); + EXPECT_EQ(interface_->Send(metadata, SerializedMessage{&byte, 1U}).GetResult(), core::Result::NotExpected); } } // namespace diff --git a/modules/jarnax/tests/mocks/jarnax/services/MockUDPSocket.hpp b/modules/jarnax/tests/mocks/jarnax/services/MockUDPDispatcher.hpp similarity index 63% rename from modules/jarnax/tests/mocks/jarnax/services/MockUDPSocket.hpp rename to modules/jarnax/tests/mocks/jarnax/services/MockUDPDispatcher.hpp index ea0bca0..e439285 100644 --- a/modules/jarnax/tests/mocks/jarnax/services/MockUDPSocket.hpp +++ b/modules/jarnax/tests/mocks/jarnax/services/MockUDPDispatcher.hpp @@ -1,25 +1,25 @@ -#ifndef JARNAX_SERVICES_MOCK_UDP_SOCKET_HPP -#define JARNAX_SERVICES_MOCK_UDP_SOCKET_HPP +#ifndef JARNAX_SERVICES_MOCK_UDP_DISPATCHER_HPP +#define JARNAX_SERVICES_MOCK_UDP_DISPATCHER_HPP #include -#include "jarnax/services/CyphalUDPSocket.hpp" +#include "jarnax/services/CyphalUDPDispatcher.hpp" namespace jarnax { namespace cyphal { namespace udp { -class MockSocket : public Socket { +class MockDispatcher : public Dispatcher { public: MOCK_METHOD(core::Status, Join, (Endpoint const& multicast_endpoint, DatagramHandler& handler), (override)); MOCK_METHOD(core::Status, Leave, (Endpoint const& multicast_endpoint), (override)); MOCK_METHOD(core::Status, Send, (Endpoint const& destination, core::Span payload), (override)); - virtual ~MockSocket() = default; + virtual ~MockDispatcher() = default; }; } // namespace udp } // namespace cyphal } // namespace jarnax -#endif // JARNAX_SERVICES_MOCK_UDP_SOCKET_HPP +#endif // JARNAX_SERVICES_MOCK_UDP_DISPATCHER_HPP \ No newline at end of file