FERS 0.1.0
The Flexible Extensible Radar Simulator
Loading...
Searching...
No Matches
serial::vita49::Vita49OutputSink Class Referencefinal

#include "vita49_output_sink.h"

+ Inheritance diagram for serial::vita49::Vita49OutputSink:
+ Collaboration diagram for serial::vita49::Vita49OutputSink:

Public Member Functions

 Vita49OutputSink (std::unique_ptr< DatagramSender > sender=nullptr, core::ReceiverOutputTelemetryCallback telemetry_callback=nullptr)
 
 ~Vita49OutputSink () override
 
 Vita49OutputSink (const Vita49OutputSink &)=delete
 
Vita49OutputSinkoperator= (const Vita49OutputSink &)=delete
 
 Vita49OutputSink (Vita49OutputSink &&)=delete
 
Vita49OutputSinkoperator= (Vita49OutputSink &&)=delete
 
void initializeRun (const core::OutputConfig &config, std::string simulation_name) override
 
std::uint32_t registerStream (const core::ReceiverStreamDescriptor &stream) override
 
void openStream (std::uint32_t stream_id, RealType first_sample_time) override
 
void submitBlock (const core::ReceiverSampleBlock &block) override
 
void submitBlocks (std::span< const core::ReceiverSampleBlock > blocks) override
 
void emitContextHeartbeat (RealType simulation_time) override
 
void closeStream (std::uint32_t stream_id) override
 
core::OutputStats finalize () override
 
core::OutputStats snapshotStats () const override
 

Detailed Description

Definition at line 25 of file vita49_output_sink.h.

Constructor & Destructor Documentation

◆ Vita49OutputSink() [1/3]

serial::vita49::Vita49OutputSink::Vita49OutputSink ( std::unique_ptr< DatagramSender sender = nullptr,
core::ReceiverOutputTelemetryCallback  telemetry_callback = nullptr 
)
explicit

Definition at line 52 of file vita49_output_sink.cpp.

53 :
54 _telemetry_callback(std::move(telemetry_callback)), _provided_sender(std::move(sender))
55 {
56 }
math::Vec3 max

◆ ~Vita49OutputSink()

serial::vita49::Vita49OutputSink::~Vita49OutputSink ( )
override

Definition at line 58 of file vita49_output_sink.cpp.

59 {
60 if (!_finalized)
61 {
62 try
63 {
64 (void)finalize();
65 }
66 catch (...)
67 {
68 }
69 }
70 }
core::OutputStats finalize() override

References finalize(), and max.

+ Here is the call graph for this function:

◆ Vita49OutputSink() [2/3]

serial::vita49::Vita49OutputSink::Vita49OutputSink ( const Vita49OutputSink )
delete

◆ Vita49OutputSink() [3/3]

serial::vita49::Vita49OutputSink::Vita49OutputSink ( Vita49OutputSink &&  )
delete

Member Function Documentation

◆ closeStream()

void serial::vita49::Vita49OutputSink::closeStream ( std::uint32_t  stream_id)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 241 of file vita49_output_sink.cpp.

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 }

References max.

◆ emitContextHeartbeat()

void serial::vita49::Vita49OutputSink::emitContextHeartbeat ( RealType  simulation_time)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 228 of file vita49_output_sink.cpp.

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 }

References max.

◆ finalize()

core::OutputStats serial::vita49::Vita49OutputSink::finalize ( )
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 254 of file vita49_output_sink.cpp.

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())
307 : _config.vita49.epoch_unix_nanoseconds,
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 }

References core::Vita49OutputConfig::epoch_unix_nanoseconds, max, core::OutputStats::mode, core::OutputConfig::vita49, and core::Vita49Udp.

Referenced by ~Vita49OutputSink().

+ Here is the caller graph for this function:

◆ initializeRun()

void serial::vita49::Vita49OutputSink::initializeRun ( const core::OutputConfig config,
std::string  simulation_name 
)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 72 of file vita49_output_sink.cpp.

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 }

References max, and core::Vita49Udp.

◆ openStream()

void serial::vita49::Vita49OutputSink::openStream ( std::uint32_t  stream_id,
RealType  first_sample_time 
)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 125 of file vita49_output_sink.cpp.

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 }

References max.

◆ operator=() [1/2]

Vita49OutputSink & serial::vita49::Vita49OutputSink::operator= ( const Vita49OutputSink )
delete

◆ operator=() [2/2]

Vita49OutputSink & serial::vita49::Vita49OutputSink::operator= ( Vita49OutputSink &&  )
delete

◆ registerStream()

std::uint32_t serial::vita49::Vita49OutputSink::registerStream ( const core::ReceiverStreamDescriptor stream)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 105 of file vita49_output_sink.cpp.

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 }
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream)

References max, core::ReceiverStreamDescriptor::mode, core::ReceiverStreamDescriptor::receiver_id, core::ReceiverStreamDescriptor::receiver_name, core::ReceiverStreamDescriptor::reference_frequency, serial::vita49::StreamRegistry::registerStream(), and core::ReceiverStreamDescriptor::sample_rate.

Referenced by submitBlocks().

+ Here is the call graph for this function:
+ Here is the caller graph for this function:

◆ snapshotStats()

core::OutputStats serial::vita49::Vita49OutputSink::snapshotStats ( ) const
overridevirtual

Reimplemented from core::ReceiverOutputSink.

Definition at line 323 of file vita49_output_sink.cpp.

324 {
325 std::scoped_lock const lock(_mutex);
326 return snapshotStatsLocked();
327 }

References max.

◆ submitBlock()

void serial::vita49::Vita49OutputSink::submitBlock ( const core::ReceiverSampleBlock block)
overridevirtual

Implements core::ReceiverOutputSink.

Definition at line 133 of file vita49_output_sink.cpp.

133{ submitBlocks({&block, 1}); }
void submitBlocks(std::span< const core::ReceiverSampleBlock > blocks) override

References max, and submitBlocks().

+ Here is the call graph for this function:

◆ submitBlocks()

void serial::vita49::Vita49OutputSink::submitBlocks ( std::span< const core::ReceiverSampleBlock blocks)
overridevirtual

Reimplemented from core::ReceiverOutputSink.

Definition at line 135 of file vita49_output_sink.cpp.

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 }
std::uint32_t registerStream(const core::ReceiverStreamDescriptor &stream) override
double RealType
Type for real numbers.
Definition config.h:27

References max, and registerStream().

Referenced by submitBlock().

+ Here is the call graph for this function:
+ Here is the caller graph for this function:

The documentation for this class was generated from the following files: