Files
MARTe-Integrated-Components/Docs/UDPStreamer.md
T
Martino FerrariandClaude Sonnet 4.6 14d5351a81 fix: UDPSClient uses two-arg Join so multicast receiver lands on the right interface
UDPSClient::ConnectMulticast was calling the single-arg BasicUDPSocket::Join,
which forwards NULL as the local interface and lets the kernel bind to
INADDR_ANY.  On a multi-homed host (or when the server sends on loopback via
Interface = "127.0.0.1") the client joins the wrong interface and silently
receives nothing.

Fix: read the optional Interface key inside the useMulticast block in
UDPSClient::Initialise; in ConnectMulticast call the two-arg
Join(group, interface) when Interface is set, and fall back to the one-arg
call otherwise to preserve the existing INADDR_ANY behaviour for configs that
omit it.

Forward the new optional Interface key through UDPStreamerClient (read from
DataSource config, written into the UDPSClient ConfigurationDatabase only
when non-empty).  Extend the "Joined multicast group" log to report the
interface name or "default".

Regression test TestExecute_MulticastReceivesDataOnInterface: mock TCP
control listener + multicast UDP DATA socket with IP_MULTICAST_IF set to
127.0.0.1, verifying a uint32 value of 424242 reaches DataSource signal
memory.  Confirmed FAILED without the Join fix and PASSED with it.

Suite: 133/133 (was 132/132 before this commit).

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
2026-08-17 08:23:41 +02:00

17 KiB
Raw Blame History

UDPStreamer DataSource

UDPStreamer is a MARTe2 output DataSource that streams signals from a real-time application to a single connected UDP client. It is fully asynchronous from the RT thread: the RT cycle only performs a fast spinlock + memcpy, while all network I/O runs on a dedicated background thread.

Key Features

  • Zero-copy RT pathSynchronise() only locks, copies signal memory, and posts a semaphore.
  • Unicast and multicast — unicast (default): single client at a time, new CONNECT replaces the previous session. Multicast: multiple clients receive data simultaneously by joining a multicast group; control traffic uses a TCP listener.
  • Packet fragmentation — large payloads are split into ≤ MaxPayloadSize-byte datagrams, each with a header carrying fragment index and total count so the client can reassemble them.
  • Signal quantizationfloat32/float64 signals can be linearly quantized to uint8, int8, uint16, or int16 on the wire, reducing bandwidth significantly.
  • Temporal arrays — signals with NumberOfElements > 1 can carry per-sample time metadata via TimeMode and TimeSignal, enabling high-frequency burst transmission (e.g. 1 000 samples per RT cycle at 1 MSps).
  • Publishing modesStrict (one packet per RT cycle), Accumulate (batch N snapshots then flush on size or time limit), Decimate (send every Nth cycle).

Configuration

+Streamer = {
    Class           = UDPStreamer

    // Network
    Port            = 44500        // UDP port (unicast) or TCP control port (multicast)
    MaxPayloadSize  = 1400         // Maximum bytes per UDP datagram (default: 1400)
                                   //   Must be > 17 (header size). Tune for MTU.

    // Multicast (optional — omit for unicast mode)
    MulticastGroup  = "239.0.0.1"  // IPv4 multicast address (224.0.0.0/4)
    Interface       = "192.168.1.10" // Local IPv4 address of the outgoing interface
                                   //   (mandatory when MulticastGroup is set)
    DataPort        = 44501        // UDP port for multicast DATA (default: Port+1)

    // Publishing mode (optional)
    PublishingMode  = "Strict"     // Strict | Accumulate | Decimate
    // For Accumulate mode:
    MinRefreshRate  = 120.0        // Flush frequency in Hz (required for Accumulate)
    // For Decimate mode:
    Ratio           = 10           // Send 1 packet every N RT cycles (required for Decimate)

    // Background thread (optional)
    CPUMask         = 0x2          // CPU affinity mask for the network thread
    StackSize       = 1048576      // Stack size in bytes (default: 1 MiB)

    Signals = {
        // ── Scalar signal ────────────────────────────────────────────────────
        Time = {
            Type = uint32
            Unit = "us"            // Optional: physical unit string (informational)
        }

        // ── Float signal with quantization ───────────────────────────────────
        Pressure = {
            Type         = float32
            Unit         = "Pa"
            RangeMin     = 0.0         // Required when QuantizedType is set
            RangeMax     = 1000000.0   // Required when QuantizedType is set
            QuantizedType = uint16     // none | uint8 | int8 | uint16 | int16
        }

        // ── Temporal array (packed burst) ────────────────────────────────────
        Channel1 = {
            Type               = float32
            NumberOfDimensions = 1
            NumberOfElements   = 1000         // N samples per RT cycle
            Unit               = "V"
            TimeMode           = FirstSample  // see Time Modes below
            TimeSignal         = Time         // name of a scalar signal in this DataSource
            SamplingRate       = 1000000.0    // Hz — used by client to reconstruct timestamps
        }
    }
}

Top-level Parameters

Parameter Type Default Description
Port uint16 44500 UDP server port (unicast) or TCP control port (multicast). Values ≤ 1024 produce a warning.
MaxPayloadSize uint32 1400 Max payload bytes per UDP datagram (min 18)
MulticastGroup string (absent) IPv4 multicast address (e.g. "239.0.0.1"). Must be in 224.0.0.0/4. Absent or empty = unicast mode.
Interface string (absent) Local IPv4 address of the interface multicast DATA leaves from, in dotted-quad form (e.g. "192.168.1.10", or "127.0.0.1" for loopback-only testing). Not an interface name — "eth0" is rejected. Mandatory when MulticastGroup is set.
DataPort uint16 Port+1 UDP port for multicast DATA datagrams. Ignored in unicast mode.
PublishingMode string Strict Strict: send every RT cycle. Accumulate: batch until size/time limit. Decimate: send every Nth cycle.
MinRefreshRate float64 Flush frequency in Hz. Required when PublishingMode = Accumulate.
Ratio uint32 Send 1 packet every Ratio RT cycles. Required when PublishingMode = Decimate.
CPUMask uint32 0xFFFFFFFF (any) Background thread CPU affinity bitmask
StackSize uint32 MARTe2 default Background thread stack size in bytes

Per-signal Parameters

Parameter Type Default Applies to
Unit string "" Any type — informational, forwarded to client in CONFIG
RangeMin float64 0.0 float32/float64 with QuantizedType
RangeMax float64 1.0 float32/float64 with QuantizedType
QuantizedType string none float32/float64 only
TimeMode string PacketTime Signals with NumberOfElements > 1
TimeSignal string Required when TimeModePacketTime
SamplingRate float64 0.0 Required when TimeMode = FirstSample or LastSample

Quantization Types

Value Wire type Bit depth Notes
none same as source Raw copy, no quantization
uint8 uint8 8-bit Maps [RangeMin, RangeMax][0, 255]
int8 int8 8-bit Maps [RangeMin, RangeMax][-127, 127]
uint16 uint16 16-bit Maps [RangeMin, RangeMax][0, 65 535]
int16 int16 16-bit Maps [RangeMin, RangeMax][-32 767, 32 767]

Quantization formula (unsigned, e.g. uint16):

normalized = clamp((value - RangeMin) / (RangeMax - RangeMin), 0.0, 1.0)
wire_value = (uint16)(normalized × 65535)

Time Modes

Value Meaning Requirements
PacketTime The HRT counter captured at Synchronise() time is used as the packet timestamp. No per-signal time metadata.
FullArray TimeSignal carries one timestamp per element (same NumberOfElements). TimeSignal must have the same NumberOfElements.
FirstSample TimeSignal is a scalar giving the timestamp of element [0]. Elements [1..N-1] are inferred at 1/SamplingRate intervals. Scalar TimeSignal; SamplingRate > 0.
LastSample Same as FirstSample but TimeSignal is the timestamp of element [N-1]. Scalar TimeSignal; SamplingRate > 0.

Network Modes

Unicast (default)

The server opens a single UDP socket on Port. The client initiates the session by sending a CONNECT packet to that port. The server replies with a CONFIG packet on the same socket and subsequently sends DATA packets directly to the client's address. One client at a time; a new CONNECT evicts the previous client.

Multicast

Enabled by setting MulticastGroup to a valid IPv4 multicast address (224.0.0.0/4). The Interface parameter is mandatory. It is the local IPv4 address of the interface DATA datagrams leave from, given in dotted-quad form — it sets IP_MULTICAST_IF on the data socket and is parsed with inet_addr(), so an interface name such as "eth0" is rejected and Initialise fails.

Receivers must join the group on the matching interface. A receiver that joins with INADDR_ANY lets the kernel pick the default-route interface, and it will silently receive nothing if that is not the interface named by Interface.

The server opens a TCP listener on Port for control traffic and a UDP socket aimed at MulticastGroup:DataPort for data traffic. The client:

  1. Connects to Port via TCP and sends a CONNECT packet.
  2. Receives the CONFIG packet over TCP.
  3. Joins the multicast group (MulticastGroup:DataPort) to receive DATA packets.

Multiple clients may receive data simultaneously by joining the same group.


Publishing Modes

Strict (default)

Sends one DATA packet for every Synchronise() call (every RT cycle). Simplest and lowest latency.

Accumulate

Batches multiple RT-cycle snapshots into a single DATA packet. All signals (scalars and arrays) are accumulated: one full snapshot per RT cycle. The batch is flushed when either:

  • Size condition: adding one more sample would exceed MaxPayloadSize.
  • Time condition: 1/MinRefreshRate seconds have elapsed since the last flush.

The maximum batch count is computed automatically from MaxPayloadSize and the total wire size of all signals. Scalar signals with Unit="us" or "ns" are auto-promoted as the per-sample FullArray time reference for all other scalars.

Requires MinRefreshRate (Hz) to be set.

Decimate

Sends one DATA packet every Ratio RT cycles, dropping intermediate cycles. Only the most recent snapshot at the Nth cycle is sent.

Requires Ratio (≥ 1) to be set. Ratio = 1 is equivalent to Strict mode (a warning is logged).


Broker

UDPStreamer uses MemoryMapSynchronisedOutputBroker for output signals. This broker is called automatically by the MARTe2 scheduler after all GAMs in the thread have executed.

Note: Input signals are not supported — GetBrokerName() returns "" for InputSignals.


Lifecycle

PrepareNextState()          ← opens UDP server socket, starts background thread
    ↓
[RT thread, each cycle]
    GAMs execute            ← write into Streamer signal memory via broker
    Synchronise()           ← spinlock + memcpy to readyBuffer + post dataSem
    ↓
[Background thread]
    Poll serverSocket       ← receive CONNECT / DISCONNECT / ACK
    Wait dataSem            ← woken by Synchronise()
    QuantizeAndSerialize()  ← build wire payload
    SendFragmented()        ← send DATA fragments to client

Performance Notes

  • The RT path (Synchronise()) performs only: FastLock() + memcpy + FastUnLock() + EventSem.Post(). No socket calls, no heap allocation.
  • readyBuffer and wireBuffer are allocated once in AllocateMemory().
  • If no client is connected the background thread skips serialisation entirely.
  • Packet loss is tolerated silently. ACK tracking is reserved for future use.

Example: minimal scalar streaming (unicast)

+Data = {
    Class = ReferenceContainer

    +Streamer = {
        Class = UDPStreamer
        Port = 44500
        Signals = {
            Counter = { Type = uint32 }
            Voltage = { Type = float32; Unit = "V"; RangeMin = -10.0; RangeMax = 10.0; QuantizedType = uint16 }
        }
    }
}

Example: multicast with accumulation

+Streamer = {
    Class            = UDPStreamer
    Port             = 44500           // TCP control port
    MulticastGroup   = "239.0.0.1"     // Enables multicast mode
    Interface        = "192.168.1.10"  // Local IP of the outgoing interface (mandatory)
    DataPort         = 44501           // UDP data port (default: Port+1)
    MaxPayloadSize   = 1400
    PublishingMode   = "Accumulate"
    MinRefreshRate   = 60.0            // Flush at least 60 times/s

    Signals = {
        Time    = { Type = uint32; Unit = "us" }
        Voltage = { Type = float32; Unit = "V"; RangeMin = -10.0; RangeMax = 10.0; QuantizedType = uint16 }
    }
}

Example: high-frequency burst

+Streamer = {
    Class          = UDPStreamer
    Port           = 44500
    MaxPayloadSize = 1400

    Signals = {
        T0 = { Type = uint32; Unit = "us" }

        Ch1 = {
            Type               = float32
            NumberOfDimensions = 1
            NumberOfElements   = 1000        // 1000 samples per RT cycle
            Unit               = "V"
            TimeMode           = FirstSample
            TimeSignal         = T0
            SamplingRate       = 1000000.0   // 1 MSps
        }
    }
}

UDPStreamerClient DataSource

UDPStreamerClient is a MARTe2 input DataSource that receives signals from a UDPStreamer server. Transport, fragment reassembly, and auto-reconnect are delegated to UDPSClient; the DataSource only decodes CONFIG/DATA payloads into real-time signal memory.

Configuration

+ClientDS = {
    Class          = UDPStreamerClient
    ServerAddress  = "192.168.1.10"  // UDPStreamer server IP
    Port           = 44500           // Server port

    // Multicast (optional — omit for unicast)
    MulticastGroup = "239.0.0.1"
    DataPort       = 44501           // UDP data port (default: Port+1)
    Interface      = "192.168.1.10"  // See table below

    MaxPayloadSize = 1400

    Signals = {
        Counter = { Type = uint32 }
    }
}

Parameters

Parameter Type Default Description
ServerAddress string 127.0.0.1 IPv4 address of the UDPStreamer server.
Port uint16 44500 Server UDP port (unicast) or TCP control port (multicast).
MulticastGroup string (absent) IPv4 multicast address. Presence enables multicast mode.
DataPort uint16 Port+1 UDP port for multicast DATA datagrams.
Interface string (absent) Local IPv4 dotted-quad address (e.g. "127.0.0.1") of the interface to join the multicast group on. Optional: omitting it uses the default-route interface (INADDR_ANY), which silently receives nothing if the server sends on a different interface. Not an interface name — "eth0" is invalid.
MaxPayloadSize uint32 1400 Max payload bytes per datagram (must match the server).
SilenceTimeout float32 1.0 Seconds of no data before auto-reconnect. 0 disables.
KeepAliveInterval uint32 15 Seconds between unicast keepalive ACKs. 0 disables.
CPUMask uint32 0xFFFFFFFF CPU affinity for the background receiver thread.
StackSize uint32 default Stack size in bytes for the receiver thread.

With MaxPayloadSize = 1400, a single 1000-element float32 signal produces:

payload = 8 B (HRT timestamp) + 4 B (T0/uint32) + 4000 B (float32×1000) = 4012 B
fragments = ceil(4012 / 1383) = 3

Example: decimated output

+Streamer = {
    Class          = UDPStreamer
    Port           = 44500
    PublishingMode = "Decimate"
    Ratio          = 10                // Send 1 packet every 10 RT cycles

    Signals = {
        Time     = { Type = uint32; Unit = "us" }
        Position = { Type = float64; Unit = "mm" }
    }
}