FERS 0.1.0
The Flexible Extensible Radar Simulator
Loading...
Searching...
No Matches
vita49_output_sink.cpp
Go to the documentation of this file.
1// SPDX-License-Identifier: GPL-2.0-only
2//
3// Copyright (c) 2026-present FERS Contributors (see AUTHORS.md).
4//
5// See the GNU GPLv2 LICENSE file in the FERS project root for more information.
6
8
9#include <algorithm>
10#include <chrono>
11#include <optional>
12#include <stdexcept>
13
14#include "core/parameters.h"
16
17namespace serial::vita49
18{
19 namespace
20 {
21 constexpr auto kStatsEmitInterval = std::chrono::milliseconds(250);
22 constexpr auto kPacketTraceEmitInterval = std::chrono::milliseconds(250);
23 constexpr std::size_t kPacketTraceBatchSize = 64;
24
25 [[nodiscard]] std::uint64_t defaultEpochNanoseconds()
26 {
27 const auto now = std::chrono::system_clock::now().time_since_epoch();
28 return static_cast<std::uint64_t>(std::chrono::duration_cast<std::chrono::nanoseconds>(now).count());
29 }
30
31 [[nodiscard]] core::Vita49Timestamp toCoreTimestamp(const Timestamp& timestamp) noexcept
32 {
33 return core::Vita49Timestamp{.integer_seconds = timestamp.integer_seconds,
34 .fractional_picoseconds = timestamp.fractional_picoseconds};
35 }
36
37 [[nodiscard]] std::optional<core::Vita49Timestamp>
38 tryCoreTimestampFromEpoch(const std::uint64_t epoch_unix_nanoseconds,
39 const RealType sample_time_seconds) noexcept
40 {
41 try
42 {
43 return toCoreTimestamp(timestampFromEpoch(epoch_unix_nanoseconds, sample_time_seconds));
44 }
45 catch (...)
46 {
47 return std::nullopt;
48 }
49 }
50 }
51
52 Vita49OutputSink::Vita49OutputSink(std::unique_ptr<DatagramSender> sender,
54 _telemetry_callback(std::move(telemetry_callback)), _provided_sender(std::move(sender))
55 {
56 }
57
59 {
60 if (!_finalized)
61 {
62 try
63 {
64 (void)finalize();
65 }
66 catch (...)
67 {
68 }
69 }
70 }
71
72 void Vita49OutputSink::initializeRun(const core::OutputConfig& config, std::string simulation_name)
73 {
74 std::scoped_lock const lock(_mutex);
76 {
77 throw std::invalid_argument("Vita49OutputSink requires VITA output mode");
78 }
79 if (config.vita49.host.empty() || config.vita49.port == 0)
80 {
81 throw std::invalid_argument("VITA output requires destination host and port");
82 }
83 if (config.vita49.queue_depth == 0)
84 {
85 throw std::invalid_argument("VITA output queue depth must be positive");
86 }
87
88 _config = config;
89 _simulation_name = std::move(simulation_name);
90 _packetizer.reset();
91 _sender = std::make_unique<PacedSender>(
92 _provided_sender ? std::move(_provided_sender) : std::make_unique<UdpSender>(), config.vita49.queue_depth);
93 _sender->open(config.vita49.host, config.vita49.port);
94 _pending_contexts.clear();
95 _pacing_started = false;
96 _last_stats_emit = std::chrono::steady_clock::time_point::min();
97 _last_packet_trace_emit = std::chrono::steady_clock::now();
98 _pending_packet_traces.clear();
99 _trace_sequence = 0;
100 _initialized = true;
101 _finalized = false;
102 emitTelemetry({}, true);
103 }
104
106 {
107 std::scoped_lock const lock(_mutex);
108 const auto stream_id = _registry.registerStream(stream);
109 if (!_streams.contains(stream_id))
110 {
111 StreamState state;
112 state.descriptor = stream;
113 state.stats.receiver_id = stream.receiver_id;
114 state.stats.receiver_name = stream.receiver_name;
115 state.stats.stream_id = stream_id;
116 state.stats.mode = stream.mode.empty() ? "unknown" : stream.mode;
117 state.stats.sample_rate = stream.sample_rate;
118 state.stats.reference_frequency = stream.reference_frequency;
119 _streams.emplace(stream_id, std::move(state));
120 emitTelemetry({}, true);
121 }
122 return stream_id;
123 }
124
125 void Vita49OutputSink::openStream(const std::uint32_t stream_id, const RealType first_sample_time)
126 {
127 std::scoped_lock const lock(_mutex);
128 emitContext(stream_id, first_sample_time, true, false);
129 stateFor(stream_id).opened = true;
130 emitTelemetry({}, true);
131 }
132
134
135 void Vita49OutputSink::submitBlocks(const std::span<const core::ReceiverSampleBlock> blocks)
136 {
137 std::scoped_lock const lock(_mutex);
138 if (!_initialized || !_sender)
139 {
140 throw std::logic_error("VITA output sink has not been initialized");
141 }
142 emitTelemetry(consumeSenderDropsLocked(), false);
143 ensurePacketizer();
144
145 struct PacketAccounting
146 {
147 std::uint32_t stream_id = 0;
148 std::uint64_t sample_count = 0;
149 RealType first_sample_time = 0.0;
150 RealType end_sample_time = 0.0;
151 core::Vita49Timestamp timestamp;
152 std::optional<core::Vita49Timestamp> end_timestamp;
153 bool over_range = false;
154 };
155
156 std::vector<SerializedPacket> packets;
157 std::vector<PacketAccounting> accounting;
158 std::vector<std::uint32_t> context_stream_ids;
159 for (const auto& block : blocks)
160 {
161 const auto stream_id = registerStream(block.stream);
162 auto& state = stateFor(stream_id);
163 if (!state.opened)
164 {
165 emitContext(stream_id, block.first_sample_time, true, false);
166 state.opened = true;
167 }
168 }
169 appendPendingContexts(packets, context_stream_ids);
170
171 for (const auto& block : blocks)
172 {
173 const auto stream_id = registerStream(block.stream);
174 auto& state = stateFor(stream_id);
175 auto result = _packetizer->packetize(block, stream_id, state.packet_counts, state.sample_loss_pending);
176 state.sample_loss_pending = false;
177 const RealType sample_rate = block.sample_rate > 0.0 ? block.sample_rate : state.stats.sample_rate;
178 for (const auto& packet : result.packets)
179 {
180 const auto end_sample_time = sample_rate > 0.0
181 ? packet.first_sample_time + static_cast<RealType>(packet.sample_count) / sample_rate
182 : packet.first_sample_time;
184 .stream_id = stream_id,
185 .sample_count = packet.sample_count,
186 .first_sample_time = packet.first_sample_time,
187 .end_sample_time = end_sample_time,
188 .timestamp = toCoreTimestamp(packet.timestamp),
189 .end_timestamp = tryCoreTimestampFromEpoch(_packetizer->epochUnixNanoseconds(), end_sample_time),
190 .over_range = packet.over_range});
191 }
192 state.stats.over_range_count += result.over_range_count;
193 packets.insert(packets.end(), std::make_move_iterator(result.packets.begin()),
194 std::make_move_iterator(result.packets.end()));
195 }
196
197 std::stable_sort(packets.begin(), packets.end(), [](const SerializedPacket& lhs, const SerializedPacket& rhs)
198 { return lhs.first_sample_time < rhs.first_sample_time; });
199 startPacing();
200 if (!enqueuePackets(std::move(packets)))
201 {
202 return;
203 }
204
205 for (const auto stream_id : context_stream_ids)
206 {
207 ++stateFor(stream_id).stats.context_packets;
208 }
209 for (const auto& packet : accounting)
210 {
211 auto& state = stateFor(packet.stream_id);
212 state.stats.samples_emitted += packet.sample_count;
213 ++state.stats.packets_emitted;
214 if (!state.stats.first_sample_time.has_value())
215 {
216 state.stats.first_sample_time = packet.first_sample_time;
217 state.stats.first_timestamp = packet.timestamp;
218 }
219 state.stats.end_sample_time = packet.end_sample_time;
220 state.stats.end_timestamp = packet.end_timestamp;
221 if (packet.over_range)
222 {
223 state.over_range_pending = true;
224 }
225 }
226 }
227
229 {
230 std::scoped_lock const lock(_mutex);
231 emitTelemetry(consumeSenderDropsLocked(), false);
232 for (auto& [stream_id, state] : _streams)
233 {
234 if (!state.closed && simulation_time - state.last_context_time >= 1.0)
235 {
236 emitContext(stream_id, simulation_time, false, false);
237 }
238 }
239 }
240
241 void Vita49OutputSink::closeStream(const std::uint32_t stream_id)
242 {
243 std::scoped_lock const lock(_mutex);
244 emitTelemetry(consumeSenderDropsLocked(), false);
245 auto& state = stateFor(stream_id);
246 if (!state.closed)
247 {
248 emitContext(stream_id, state.last_context_time, false, true);
249 state.closed = true;
250 emitTelemetry({}, true);
251 }
252 }
253
255 {
256 std::scoped_lock const lock(_mutex);
257 if (_finalized)
258 {
259 auto stats = snapshotStatsLocked();
260 emitTelemetry({}, true);
261 return stats;
262 }
263
264 if (_sender)
265 {
266 _sender->flush();
267 emitTelemetry(consumeSenderDropsLocked(), false);
268 }
269
270 for (auto& [stream_id, state] : _streams)
271 {
272 if (!state.closed)
273 {
274 emitContext(stream_id, state.last_context_time, false, true);
275 state.closed = true;
276 }
277 }
278
279 if (!_pacing_started && !_pending_contexts.empty())
280 {
281 ensurePacketizer();
282 std::vector<SerializedPacket> packets;
283 std::vector<std::uint32_t> context_stream_ids;
284 appendPendingContexts(packets, context_stream_ids);
285 std::stable_sort(packets.begin(), packets.end(),
286 [](const SerializedPacket& lhs, const SerializedPacket& rhs)
287 { return lhs.first_sample_time < rhs.first_sample_time; });
288 startPacing();
289 if (enqueuePackets(std::move(packets)))
290 {
291 for (const auto stream_id : context_stream_ids)
292 {
293 ++stateFor(stream_id).stats.context_packets;
294 }
295 }
296 }
297
298 if (_sender)
299 {
300 _sender->stop();
301 emitTelemetry(consumeSenderDropsLocked(), false);
302 }
303
305 .epoch_unix_nanoseconds = _packetizer
306 ? std::optional<std::uint64_t>(_packetizer->epochUnixNanoseconds())
308 .streams = {}};
309 for (auto& [stream_id, state] : _streams)
310 {
311 if (_sender)
312 {
313 state.stats.late_data_packet_count = _sender->lateDataPacketCount(stream_id);
314 state.stats.late_context_packet_count = _sender->lateContextPacketCount(stream_id);
315 }
316 stats.streams.push_back(state.stats);
317 }
318 _finalized = true;
319 emitTelemetry({}, true);
320 return stats;
321 }
322
324 {
325 std::scoped_lock const lock(_mutex);
326 return snapshotStatsLocked();
327 }
328
329 Vita49OutputSink::StreamState& Vita49OutputSink::stateFor(const std::uint32_t stream_id)
330 {
331 const auto found = _streams.find(stream_id);
332 if (found == _streams.end())
333 {
334 throw std::out_of_range("Unknown VITA stream ID");
335 }
336 return found->second;
337 }
338
339 const Vita49OutputSink::StreamState& Vita49OutputSink::stateFor(const std::uint32_t stream_id) const
340 {
341 const auto found = _streams.find(stream_id);
342 if (found == _streams.end())
343 {
344 throw std::out_of_range("Unknown VITA stream ID");
345 }
346 return found->second;
347 }
348
349 void Vita49OutputSink::ensurePacketizer()
350 {
351 if (_packetizer)
352 {
353 return;
354 }
355 const auto epoch_ns = _config.vita49.epoch_unix_nanoseconds.value_or(defaultEpochNanoseconds());
356 _packetizer =
357 std::make_unique<Vita49Packetizer>(epoch_ns, _config.vita49.adc_fullscale, _config.vita49.max_udp_payload);
358 }
359
360 void Vita49OutputSink::startPacing()
361 {
362 if (_pacing_started)
363 {
364 return;
365 }
366 if (!_sender)
367 {
368 throw std::logic_error("VITA paced sender is unavailable");
369 }
370 _sender->start(params::startTime());
371 _pacing_started = true;
372 }
373
374 void Vita49OutputSink::appendPendingContexts(std::vector<SerializedPacket>& packets,
375 std::vector<std::uint32_t>& context_stream_ids)
376 {
377 packets.reserve(packets.size() + _pending_contexts.size());
378 context_stream_ids.reserve(context_stream_ids.size() + _pending_contexts.size());
379 for (const auto& pending : _pending_contexts)
380 {
381 packets.push_back(buildContextPacket(pending.stream_id, pending.simulation_time, pending.stream_open,
382 pending.stream_close));
383 context_stream_ids.push_back(pending.stream_id);
384 }
385 _pending_contexts.clear();
386 }
387
388 core::OutputStats Vita49OutputSink::snapshotStatsLocked() const
389 {
391 .epoch_unix_nanoseconds = _packetizer
392 ? std::optional<std::uint64_t>(_packetizer->epochUnixNanoseconds())
393 : _config.vita49.epoch_unix_nanoseconds,
394 .streams = {}};
395 for (const auto& [stream_id, state] : _streams)
396 {
397 auto stream_stats = state.stats;
398 if (_sender)
399 {
400 stream_stats.late_data_packet_count = _sender->lateDataPacketCount(stream_id);
401 stream_stats.late_context_packet_count = _sender->lateContextPacketCount(stream_id);
402 }
403 stats.streams.push_back(std::move(stream_stats));
404 }
405 return stats;
406 }
407
408 std::vector<core::ReceiverOutputPacketTrace> Vita49OutputSink::consumeSenderDropsLocked()
409 {
410 std::vector<core::ReceiverOutputPacketTrace> traces;
411 if (!_sender)
412 {
413 return traces;
414 }
415
416 for (const auto& dropped : _sender->consumeDroppedDatagrams())
417 {
418 applyDropped(dropped);
419 if (_config.vita49.packet_trace_enabled)
420 {
421 traces.push_back(makeDropTrace(dropped));
422 }
423 }
424 return traces;
425 }
426
427 bool Vita49OutputSink::enqueuePacket(SerializedPacket&& packet)
428 {
429 if (!_sender)
430 {
431 throw std::logic_error("VITA paced sender is unavailable");
432 }
433
434 std::vector<core::ReceiverOutputPacketTrace> traces;
435 std::optional<core::ReceiverOutputPacketTrace> sent_trace;
436 if (_config.vita49.packet_trace_enabled)
437 {
438 sent_trace = makeTrace(packet, packet.context_packet ? "context" : "data");
439 }
440 const auto result = _sender->enqueue(std::move(packet));
441 if (_config.vita49.packet_trace_enabled && result.dropped)
442 {
443 traces.push_back(makeDropTrace(*result.dropped));
444 }
445 if (result.dropped)
446 {
447 applyDropped(*result.dropped);
448 }
449 if (result.enqueued && sent_trace.has_value())
450 {
451 traces.push_back(std::move(*sent_trace));
452 }
453 emitTelemetry(std::move(traces), false);
454 return result.enqueued;
455 }
456
457 bool Vita49OutputSink::enqueuePackets(std::vector<SerializedPacket> packets)
458 {
459 if (!_sender)
460 {
461 throw std::logic_error("VITA paced sender is unavailable");
462 }
463
464 std::vector<core::ReceiverOutputPacketTrace> traces;
465 if (_config.vita49.packet_trace_enabled)
466 {
467 traces.reserve(packets.size());
468 for (const auto& packet : packets)
469 {
470 traces.push_back(makeTrace(packet, packet.context_packet ? "context" : "data"));
471 }
472 }
473 const bool enqueued = _sender->enqueueBatch(std::move(packets));
474 if (!enqueued)
475 {
476 traces.clear();
477 }
478 emitTelemetry(std::move(traces), false);
479 return enqueued;
480 }
481
482 void Vita49OutputSink::emitTelemetry(std::vector<core::ReceiverOutputPacketTrace> packets, const bool force_stats)
483 {
484 if (!_telemetry_callback)
485 {
486 return;
487 }
488
489 std::optional<core::OutputStats> stats;
490 const auto now = std::chrono::steady_clock::now();
491 if (force_stats || _last_stats_emit == std::chrono::steady_clock::time_point::min() ||
492 now - _last_stats_emit >= kStatsEmitInterval)
493 {
494 stats = snapshotStatsLocked();
495 _last_stats_emit = now;
496 }
497
498 for (auto& packet : packets)
499 {
500 packet.sequence = ++_trace_sequence;
501 _pending_packet_traces.push_back(std::move(packet));
502 }
503
504 std::vector<core::ReceiverOutputPacketTrace> packet_batch;
505 const bool trace_interval_elapsed = _last_packet_trace_emit == std::chrono::steady_clock::time_point::min() ||
506 now - _last_packet_trace_emit >= kPacketTraceEmitInterval;
507 if (!_pending_packet_traces.empty() &&
508 (force_stats || _pending_packet_traces.size() >= kPacketTraceBatchSize || trace_interval_elapsed))
509 {
510 packet_batch = std::move(_pending_packet_traces);
511 _pending_packet_traces.clear();
512 _last_packet_trace_emit = now;
513 }
514
515 if (!stats.has_value() && packet_batch.empty())
516 {
517 return;
518 }
519 _telemetry_callback(stats, packet_batch);
520 }
521
522 core::ReceiverOutputPacketTrace Vita49OutputSink::makeTrace(const SerializedPacket& packet, std::string event) const
523 {
525 .event = std::move(event),
526 .stream_id = packet.stream_id,
527 .byte_count = packet.bytes.size(),
528 .sample_count = packet.sample_count,
529 .first_sample_time = packet.first_sample_time,
530 .timestamp = toCoreTimestamp(packet.timestamp),
531 .data_packet = packet.data_packet,
532 .context_packet = packet.context_packet,
533 .dropped = false,
534 .over_range = packet.over_range,
535 .sample_loss = packet.sample_loss};
536 }
537
538 core::ReceiverOutputPacketTrace Vita49OutputSink::makeDropTrace(const DroppedDatagram& dropped) const
539 {
541 .event = "drop",
542 .stream_id = dropped.stream_id,
543 .byte_count = 0,
544 .sample_count = dropped.sample_count,
545 .first_sample_time = 0.0,
546 .timestamp = std::nullopt,
547 .data_packet = dropped.data_packet,
548 .context_packet = dropped.context_packet,
549 .dropped = true,
550 .over_range = false,
551 .sample_loss = true};
552 }
553
554 SerializedPacket Vita49OutputSink::buildContextPacket(const std::uint32_t stream_id, const RealType simulation_time,
555 const bool stream_open, const bool stream_close)
556 {
557 if (!_packetizer)
558 {
559 throw std::logic_error("VITA packetizer is unavailable");
560 }
561 auto& state = stateFor(stream_id);
562 const RealType context_time = simulation_time <= -1.0e200 ? 0.0 : simulation_time;
563 const auto timestamp = timestampFromEpoch(_packetizer->epochUnixNanoseconds(), context_time);
564 const ContextBuildRequest request{.stream = state.descriptor,
565 .stream_id = stream_id,
566 .simulation_name = _simulation_name,
567 .adc_fullscale = _packetizer->adcFullscale(),
568 .timestamp = timestamp,
569 .packet_count = state.packet_counts.next(),
570 .valid_data = true,
571 .calibrated_time = true,
572 .reference_lock = true,
573 .over_range = state.over_range_pending,
574 .sample_loss = state.sample_loss_pending,
575 .stream_open = stream_open,
576 .stream_close = stream_close};
577 const auto context = Vita49ContextBuilder::build(request);
578 auto packet = _packetizer->makeContextPacket(context);
579 packet.first_sample_time = context_time;
580 state.last_context_time = context_time;
581 state.sample_loss_pending = false;
582 state.over_range_pending = false;
583 return packet;
584 }
585
586 void Vita49OutputSink::emitContext(const std::uint32_t stream_id, const RealType simulation_time,
587 const bool stream_open, const bool stream_close)
588 {
589 if (!_pacing_started)
590 {
591 const RealType context_time = simulation_time <= -1.0e200 ? 0.0 : simulation_time;
592 _pending_contexts.push_back(PendingContext{.stream_id = stream_id,
593 .simulation_time = context_time,
594 .stream_open = stream_open,
595 .stream_close = stream_close});
596 stateFor(stream_id).last_context_time = context_time;
597 return;
598 }
599 auto packet = buildContextPacket(stream_id, simulation_time, stream_open, stream_close);
600 if (enqueuePacket(std::move(packet)))
601 {
602 ++stateFor(stream_id).stats.context_packets;
603 }
604 }
605
606 void Vita49OutputSink::applyDropped(const DroppedDatagram& dropped)
607 {
608 if (dropped.stream_id == 0 || !_streams.contains(dropped.stream_id))
609 {
610 return;
611 }
612 auto& state = stateFor(dropped.stream_id);
613 ++state.stats.packets_dropped;
614 state.stats.samples_dropped += dropped.sample_count;
615 if (dropped.data_packet || (!dropped.context_packet && dropped.sample_count > 0))
616 {
617 state.stats.packets_emitted -= std::min<std::uint64_t>(state.stats.packets_emitted, 1u);
618 state.stats.samples_emitted -= std::min(state.stats.samples_emitted, dropped.sample_count);
619 }
620 else if (dropped.context_packet)
621 {
622 state.stats.context_packets -= std::min<std::uint64_t>(state.stats.context_packets, 1u);
623 }
624 state.sample_loss_pending = true;
625 }
626
627 std::unique_ptr<core::ReceiverOutputSink>
629 {
630 return std::make_unique<Vita49OutputSink>(nullptr, std::move(telemetry_callback));
631 }
632}
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream)
void closeStream(std::uint32_t stream_id) override
core::OutputStats snapshotStats() const override
void emitContextHeartbeat(RealType simulation_time) override
void submitBlocks(std::span< const core::ReceiverSampleBlock > blocks) override
core::OutputStats finalize() override
void openStream(std::uint32_t stream_id, RealType first_sample_time) override
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream) override
void submitBlock(const core::ReceiverSampleBlock &block) override
void initializeRun(const core::OutputConfig &config, std::string simulation_name) override
Vita49OutputSink(std::unique_ptr< DatagramSender > sender=nullptr, core::ReceiverOutputTelemetryCallback telemetry_callback=nullptr)
double RealType
Type for real numbers.
Definition config.h:27
std::function< void(const std::optional< OutputStats > &, std::span< const ReceiverOutputPacketTrace >)> ReceiverOutputTelemetryCallback
RealType startTime() noexcept
Get the start time for the simulation.
Definition parameters.h:103
Timestamp timestampFromEpoch(const std::uint64_t epoch_unix_nanoseconds, const RealType sample_time_seconds)
std::unique_ptr< core::ReceiverOutputSink > makeVita49OutputSink(core::ReceiverOutputTelemetryCallback telemetry_callback)
Defines the Parameters struct and provides methods for managing simulation parameters.
math::Vec3 max
Vita49OutputConfig vita49
std::uint16_t max_udp_payload
std::optional< std::uint64_t > epoch_unix_nanoseconds
std::uint32_t integer_seconds