15#if defined(_WIN32) || defined(_WIN64)
19#ifndef WIN32_LEAN_AND_MEAN
20#define WIN32_LEAN_AND_MEAN
25#pragma comment(lib, "ws2_32.lib")
26#define CORIUM_HAS_UDP_SOCKETS 1
27#elif __has_include(<sys/socket.h>) && __has_include(<netinet/in.h>) && __has_include(<arpa/inet.h>) && __has_include(<unistd.h>) && __has_include(<fcntl.h>)
30#include <netinet/in.h>
31#include <sys/socket.h>
33#define CORIUM_HAS_UDP_SOCKETS 1
35#define CORIUM_HAS_UDP_SOCKETS 0
47template <
size_t MaxPacketSize = 512>
77 bool openAndBind(uint16_t port = 0,
const char* ip =
"0.0.0.0") noexcept {
78#if CORIUM_HAS_UDP_SOCKETS
81#if defined(_WIN32) || defined(_WIN64)
83 WSAStartup(MAKEWORD(2, 2), &wsaData);
84 SOCKET sock = ::socket(AF_INET, SOCK_DGRAM, IPPROTO_UDP);
85 if (sock == INVALID_SOCKET) {
88 m_fd =
static_cast<int>(sock);
90 m_fd = ::socket(AF_INET, SOCK_DGRAM, 0);
97 addr.sin_family = AF_INET;
98 addr.sin_port = htons(port);
99#if defined(_WIN32) || defined(_WIN64)
100 InetPtonA(AF_INET, ip, &addr.sin_addr);
102 inet_pton(AF_INET, ip, &addr.sin_addr);
105 if (::bind(m_fd,
reinterpret_cast<const sockaddr*
>(&addr),
sizeof(addr)) < 0) {
121#if CORIUM_HAS_UDP_SOCKETS
125#if defined(_WIN32) || defined(_WIN64)
126 u_long mode = nonBlocking ? 1 : 0;
127 return ioctlsocket(
static_cast<SOCKET
>(m_fd), FIONBIO, &mode) == 0;
129 int flags = fcntl(m_fd, F_GETFL, 0);
130 if (flags < 0)
return false;
131 flags = nonBlocking ? (flags | O_NONBLOCK) : (flags & ~O_NONBLOCK);
132 return fcntl(m_fd, F_SETFL, flags) == 0;
145 bool sendTo(
const char* ip, uint16_t port, std::span<const uint8_t> data)
noexcept {
146#if CORIUM_HAS_UDP_SOCKETS
154 sockaddr_in destAddr{};
155 destAddr.sin_family = AF_INET;
156 destAddr.sin_port = htons(port);
157#if defined(_WIN32) || defined(_WIN64)
158 InetPtonA(AF_INET, ip, &destAddr.sin_addr);
160 static_cast<SOCKET
>(m_fd),
161 reinterpret_cast<const char*
>(data.data()),
162 static_cast<int>(data.size()),
164 reinterpret_cast<const sockaddr*
>(&destAddr),
167 return sent > 0 &&
static_cast<size_t>(sent) == data.size();
169 inet_pton(AF_INET, ip, &destAddr.sin_addr);
170 auto sent = ::sendto(
175 reinterpret_cast<const sockaddr*
>(&destAddr),
178 return sent > 0 &&
static_cast<size_t>(sent) == data.size();
195 template <
typename Event,
typename EventVariant>
197 auto packet = corium::wire::WireSerializer::serialize<Event, EventVariant, MaxPacketSize>(event);
198 return sendTo(ip, port, std::span<const uint8_t>(
199 reinterpret_cast<const uint8_t*
>(&packet), packet.totalWireSize()));
206 bool receive(std::span<uint8_t> bufferOut,
size_t& bytesReceivedOut)
noexcept {
207#if CORIUM_HAS_UDP_SOCKETS
209 bytesReceivedOut = 0;
213#if defined(_WIN32) || defined(_WIN64)
214 int recvd = ::recvfrom(
215 static_cast<SOCKET
>(m_fd),
216 reinterpret_cast<char*
>(bufferOut.data()),
217 static_cast<int>(bufferOut.size()),
221 auto recvd = ::recvfrom(
229 bytesReceivedOut = 0;
232 bytesReceivedOut =
static_cast<size_t>(recvd);
236 bytesReceivedOut = 0;
247 template <
typename EventVariant,
typename Sink>
250 if (!
receive(std::span<uint8_t>(m_rxBuffer.data(), m_rxBuffer.size()), recvd)) {
259 std::memcpy(&packet, m_rxBuffer.data(), recvd >
sizeof(packet) ?
sizeof(packet) : recvd);
260 return corium::wire::WireSerializer::deserializeAndPush<EventVariant, MaxPacketSize, Sink>(
261 packet, sink, priority);
266#if CORIUM_HAS_UDP_SOCKETS
268#if defined(_WIN32) || defined(_WIN64)
269 closesocket(
static_cast<SOCKET
>(m_fd));
279 [[nodiscard]]
bool isOpen() const noexcept {
Bounded and multi-tier priority MPSC queueing policies.
Type-safe serialization and direct event sink deserialization.
Binary packet framing with CRC-16 checksum and schema versioning.
Statically buffered, zero-heap UDP communication channel for distributed Corium nodes.
Definition StaticUdpChannel.hpp:48
StaticUdpChannel & operator=(StaticUdpChannel &&other) noexcept
Definition StaticUdpChannel.hpp:64
StaticUdpChannel & operator=(const StaticUdpChannel &)=delete
bool sendTo(const char *ip, uint16_t port, std::span< const uint8_t > data) noexcept
Send raw payload to target UDP endpoint.
Definition StaticUdpChannel.hpp:145
StaticUdpChannel(StaticUdpChannel &&other) noexcept
Definition StaticUdpChannel.hpp:59
bool setNonBlocking(bool nonBlocking=true) noexcept
Set socket non-blocking mode.
Definition StaticUdpChannel.hpp:120
bool openAndBind(uint16_t port=0, const char *ip="0.0.0.0") noexcept
Open UDP socket and bind to a local port and IP address.
Definition StaticUdpChannel.hpp:77
bool receive(std::span< uint8_t > bufferOut, size_t &bytesReceivedOut) noexcept
Receive raw bytes from incoming UDP datagram.
Definition StaticUdpChannel.hpp:206
bool sendEvent(const char *ip, uint16_t port, const Event &event) noexcept
Serialize and send a typed event over UDP using Corium WirePacket framing.
Definition StaticUdpChannel.hpp:196
int nativeHandle() const noexcept
Native socket file descriptor or handle.
Definition StaticUdpChannel.hpp:284
bool isOpen() const noexcept
Returns true if socket is open and bound.
Definition StaticUdpChannel.hpp:279
bool receiveAndPush(Sink &sink, EventPriority priority=EventPriority::Normal) noexcept
Receive a WirePacket and deserialize directly into a Corium EventSink.
Definition StaticUdpChannel.hpp:248
constexpr StaticUdpChannel() noexcept=default
StaticUdpChannel(const StaticUdpChannel &)=delete
void close() noexcept
Close underlying socket.
Definition StaticUdpChannel.hpp:265
Definition StaticUdpChannel.hpp:42
EventPriority
Event priority levels for multi-priority queue policies.
Definition QueuePolicies.hpp:32
DefaultEvents Event
Alias for DefaultEvents.
Definition Events.hpp:65
Statically-sized zero-heap binary wire packet.
Definition WirePacket.hpp:64