Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions applications/nucleo-cyphal/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
52 changes: 52 additions & 0 deletions applications/nucleo-cyphal/include/HyphaUdpDispatcher.hpp
Original file line number Diff line number Diff line change
@@ -0,0 +1,52 @@
#ifndef APP_HYPHA_UDP_DISPATCHER_HPP
#define APP_HYPHA_UDP_DISPATCHER_HPP

#include <cstddef>

#include <core/Array.hpp>

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<std::uint8_t const> 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<Registration, MaxEndpoints> registrations_{};
};

} // namespace cyphal
} // namespace nucleo

#endif // APP_HYPHA_UDP_DISPATCHER_HPP
96 changes: 96 additions & 0 deletions applications/nucleo-cyphal/source/HyphaUdpDispatcher.cpp
Original file line number Diff line number Diff line change
@@ -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<std::uint8_t>(address >> 24U), static_cast<std::uint8_t>(address >> 16U),
static_cast<std::uint8_t>(address >> 8U), static_cast<std::uint8_t>(address)};
}

std::uint32_t HyphaUdpDispatcher::ToEndpointAddress(HyphaIpIPv4Address_t address) {
return (static_cast<std::uint32_t>(address.a) << 24U) | (static_cast<std::uint32_t>(address.b) << 16U) |
(static_cast<std::uint32_t>(address.c) << 8U) | static_cast<std::uint32_t>(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 = &registration;
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<std::uint8_t const> 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<std::uint8_t*>(payload.data()), static_cast<std::uint32_t>(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<std::uint8_t*>(datagram.pointer), HyphaIpSpanSize(datagram));
return;
}
}
}

} // namespace cyphal
} // namespace nucleo
Original file line number Diff line number Diff line change
@@ -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 <cstddef>
#include <cstdint>
Expand All @@ -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:
Expand All @@ -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<std::uint8_t const> 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
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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.
Expand All @@ -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 ===
Expand Down Expand Up @@ -118,15 +118,15 @@ 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_;

UdpardMemoryResource tx_memory_;
UdpardRxMemoryResources rx_memory_;
UdpardTx tx_;
UdpardRxRPCDispatcher dispatcher_;
UdpardRxRPCDispatcher rpc_dispatcher_;
udp::Endpoint service_endpoint_{};
bool service_group_joined_{false};

Expand Down
Loading
Loading